ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

使用 Airflow 构建端到端数据管道:从 CSV 下载到 Postgres 清洗入库的完整实战

使用 Airflow 构建端到端数据管道:从 CSV 下载到 Postgres 清洗入库的完整实战 使用 Airflow 构建端到端数据管道从 CSV 下载到 Postgres 清洗入库的完整实战【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文是 Apache Airflow 入门系列教程的第三篇将带你从零搭建一条小而完整的真实数据管道从外部数据源下载 CSV 文件、加载到 Postgres 暂存表staging table、清洗去重后合并进目标表。读完本文你将掌握SQLExecuteQueryOperator与PostgresHook的组合用法、Airflow UI 中 Connection连接的配置方式、以及task与dag装饰器下的 Dag 编排模式并能在本地 Docker 环境中亲手运行出完整的 ETL 流程。本教程对应的原始文档为 airflow-core/docs/tutorial/pipeline.rst配套示例数据为 pipeline_example.csv。教程概览与学习目标在此之前你应该已经写过第一个 Dag 并使用过一些基础 Operator。本教程引入的SQLExecuteQueryOperator是 Airflow 中执行 SQL 的现代、灵活方式它隶属于apache-airflow-providers-common-sql提供方包位于 providers/common/sql/src/airflow/providers/common/sql/operators/sql.py。我们将用它连接一个在 Airflow UI 中配置的本地 Postgres 数据库最终完成一条具备以下能力的管道下载一个 CSV 文件将数据加载到暂存表staging table清洗数据并通过 upsert插入或更新写入目标表在这个过程中你会获得对 Airflow UI、连接Connection系统、SQL 执行以及 Dag 编写模式的实战经验。初始环境搭建用 Docker Compose 启动 Airflow运行本教程需要安装 Docker我们将使用 Docker Compose 在本地拉起 Airflow避免对系统环境做任何全局安装既简单又安全。打开终端执行# 下载 docker-compose.yaml 文件 curl -LfO https://airflow.apache.org/docs/apache-airflow/stable/docker-compose.yaml # 创建预期的目录并设置预期的环境变量 mkdir -p ./dags ./logs ./plugins echo -e AIRFLOW_UID$(id -u) .env # 初始化数据库 docker compose up airflow-init # 启动所有服务 docker compose up注意AIRFLOW_UID$(id -u)将当前用户 ID 写入.env确保容器内的文件挂载权限与宿主机一致这是 Docker 环境下最常见的权限问题规避手段。Airflow 启动后访问http://localhost:8080进入 UI使用以下凭据登录用户名airflow密码airflow登录后你将进入 Airflow 仪表盘Dashboard在这里可以触发 Dag、查看日志以及管理整个环境。dags/目录正是后续保存 Dag 文件的位置教程中的 Dag 应保存为dags/process_employees.py。创建 Postgres 连接Connection在管道写入 Postgres 之前需要先告诉 Airflow 如何连接到数据库。在 UI 中打开Admin Connections页面点击按钮新建一条连接填写以下信息字段值Connection IDtutorial_pg_connConnection TypepostgresHostpostgresDatabaseairflow容器中的默认数据库LoginairflowPasswordairflowPort5432保存后Airflow 就掌握了如何访问 Docker 环境中运行的 Postgres 数据库。连接系统是 Airflow 的核心抽象SQLExecuteQueryOperator通过conn_idtutorial_pg_conn引用这条连接而PostgresHook则通过postgres_conn_idtutorial_pg_conn复用同一份凭据信息。从源码看连接类型postgres与PostgresHook的conn_type postgres、default_conn_name postgres_default一一对应见 providers/postgres/src/airflow/providers/postgres/hooks/postgres.py意味着你可以在连接中设置sslmode、cursor等额外参数extra 字段例如{sslmode: require}。创建暂存表与目标表数据管道的第一步是建表。我们将创建两张表employees_temp暂存表staging table存放原始数据employees清洗去重后的目标表destination使用SQLExecuteQueryOperator执行建表 SQLfrom airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator create_employees_table SQLExecuteQueryOperator( task_idcreate_employees_table, conn_idtutorial_pg_conn, sql CREATE TABLE IF NOT EXISTS employees ( Serial Number NUMERIC PRIMARY KEY, Company Name TEXT, Employee Markme TEXT, Description TEXT, Leave INTEGER );, ) create_employees_temp_table SQLExecuteQueryOperator( task_idcreate_employees_temp_table, conn_idtutorial_pg_conn, sql DROP TABLE IF EXISTS employees_temp; CREATE TABLE employees_temp ( Serial Number NUMERIC PRIMARY KEY, Company Name TEXT, Employee Markme TEXT, Description TEXT, Leave INTEGER );, )两个 task 都通过conn_idtutorial_pg_conn定位到刚刚配置的 Postgres 连接。employees_temp使用DROP TABLE IF EXISTS保证每次运行都从干净状态开始这是暂存表的标准做法employees使用CREATE TABLE IF NOT EXISTS保留历史数据为后续 upsert 做准备。注意列名带空格因此 SQL 中必须使用双引号包裹。进阶技巧sql参数不仅支持内联字符串还支持指向.sql文件的路径文件必须带.sql扩展名。从源码可见SQLExecuteQueryOperator的template_fields (sql, parameters, ...)且template_ext (.sql, .json)见 sql.py这意味着sql会被 Airflow 模板引擎渲染你可以把 SQL 语句放在dags/目录下的.sql文件中通过文件路径引用让 Dag 代码保持干净整洁还能利用 Jinja 模板变量动态生成 SQL。SQLExecuteQueryOperator 核心参数速查结合源码sql.py该 Operator 支持的关键参数包括参数默认值说明sql必填要执行的 SQL 代码或指向模板文件的路径支持.sql/.json扩展名会被模板渲染conn_idNone连接 ID决定目标数据库autocommitFalse为True时每条命令自动提交parametersNone渲染 SQL 查询时使用的参数字典或可迭代对象handlerfetch_all_handler作用于 cursor 的结果处理函数split_statementsNone是否将单个 SQL 字符串按语句拆分为None时沿用底层 hook 的run方法默认值return_lastTrue仅返回最后一条语句的结果show_return_value_in_logsFalse为True时将 Operator 输出打印到任务日志不建议对大结果集开启requires_result_fetchFalse为True时确保在完成执行前获取查询结果从execute方法的实现看sql.py该 Operator 最终调用的是数据库 hook 的run()方法当do_xcom_push为True时查询结果会自动写入 XCom供下游任务消费。下载 CSV 并加载到暂存表接下来下载 CSV 文件、保存到本地并使用PostgresHook将其加载进employees_tempimport os import requests from airflow.sdk import task from airflow.providers.postgres.hooks.postgres import PostgresHook task def get_data(): # NOTE: 请根据你的 Airflow 环境调整此路径 data_path /opt/airflow/dags/files/employees.csv os.makedirs(os.path.dirname(data_path), exist_okTrue) url https://raw.githubusercontent.com/apache/airflow/main/airflow-core/docs/tutorial/pipeline_example.csv response requests.request(GET, url) with open(data_path, w) as file: file.write(response.text) postgres_hook PostgresHook(postgres_conn_idtutorial_pg_conn) conn postgres_hook.get_conn() cur conn.cursor() with open(data_path, r) as file: cur.copy_expert( COPY employees_temp FROM STDIN WITH CSV HEADER DELIMITER AS , QUOTE \, file, ) conn.commit()这段代码展示了真实世界管道中的常见模式把 Airflow 的 task 调度能力与原生 Python 逻辑、数据库 hook 结合起来。要点拆解data_path指向 Docker 容器内/opt/airflow/dags/files/employees.csv并用os.makedirs(..., exist_okTrue)确保目录存在使用requests从示例 URL 下载 CSV本仓库中对应的示例数据为 pipeline_example.csv首行为Serial Number,Company Name,Employee Markme,Description,Leave表头创建PostgresHook获取数据库连接与游标调用cur.copy_expert(...)执行 Postgres 原生COPY ... FROM STDIN快速导入CSV HEADER表示跳过首行表头DELIMITER AS ,指定逗号分隔符QUOTE 指定双引号作为引用符——这是批量灌入数据的最高效方式之一远快于逐行INSERT显式conn.commit()提交事务。copy_expert是PostgresHook提供的专有方法见 postgres.py它将 SQL 与文件流直接交给底层 psycopg 驱动执行非常适合这种文件到表的加载场景。合并并清洗数据数据进入暂存表后需要去重并合并到最终表。我们编写一个 task 执行 SQLINSERT ... ON CONFLICT DO UPDATEfrom airflow.sdk import task from airflow.providers.postgres.hooks.postgres import PostgresHook task def merge_data(): query INSERT INTO employees SELECT * FROM ( SELECT DISTINCT * FROM employees_temp ) t ON CONFLICT (Serial Number) DO UPDATE SET Employee Markme excluded.Employee Markme, Description excluded.Description, Leave excluded.Leave; try: postgres_hook PostgresHook(postgres_conn_idtutorial_pg_conn) conn postgres_hook.get_conn() cur conn.cursor() cur.execute(query) conn.commit() return 0 except Exception as e: return 1这是 PostgreSQL 实现upsert存在则更新、不存在则插入的经典写法也是真实数据管道中的核心清洗逻辑内层SELECT DISTINCT * FROM employees_temp先对暂存数据按全字段去重外层INSERT INTO employees SELECT * ...将去重结果写入目标表ON CONFLICT (Serial Number) DO UPDATE以主键Serial Number为冲突判定依据冲突时更新Employee Markme、Description、Leave三个字段通过excluded.前缀引用本次插入尝试的值从而保证同一记录只保留最新、最干净的一份。函数返回0/1作为成功/失败信号这种显式返回状态值的写法便于后续在 XCom 中传递结果或触发告警。组装 Dag完整代码现在把所有 task 组装成一个完整的 Dag。将下面代码保存为dags/process_employees.pyimport datetime import pendulum import os import requests from airflow.sdk import dag, task from airflow.providers.postgres.hooks.postgres import PostgresHook from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator dag( dag_idprocess_employees, schedule0 0 * * *, start_datependulum.datetime(2021, 1, 1, tzUTC), catchupFalse, dagrun_timeoutdatetime.timedelta(minutes60), ) def ProcessEmployees(): create_employees_table SQLExecuteQueryOperator( task_idcreate_employees_table, conn_idtutorial_pg_conn, sql CREATE TABLE IF NOT EXISTS employees ( Serial Number NUMERIC PRIMARY KEY, Company Name TEXT, Employee Markme TEXT, Description TEXT, Leave INTEGER );, ) create_employees_temp_table SQLExecuteQueryOperator( task_idcreate_employees_temp_table, conn_idtutorial_pg_conn, sql DROP TABLE IF EXISTS employees_temp; CREATE TABLE employees_temp ( Serial Number NUMERIC PRIMARY KEY, Company Name TEXT, Employee Markme TEXT, Description TEXT, Leave INTEGER );, ) task def get_data(): # NOTE: 请根据你的 Airflow 环境调整此路径 data_path /opt/airflow/dags/files/employees.csv os.makedirs(os.path.dirname(data_path), exist_okTrue) url https://raw.githubusercontent.com/apache/airflow/main/airflow-core/docs/tutorial/pipeline_example.csv response requests.request(GET, url) with open(data_path, w) as file: file.write(response.text) postgres_hook PostgresHook(postgres_conn_idtutorial_pg_conn) conn postgres_hook.get_conn() cur conn.cursor() with open(data_path, r) as file: cur.copy_expert( COPY employees_temp FROM STDIN WITH CSV HEADER DELIMITER AS , QUOTE \, file, ) conn.commit() task def merge_data(): query INSERT INTO employees SELECT * FROM ( SELECT DISTINCT * FROM employees_temp ) t ON CONFLICT (Serial Number) DO UPDATE SET Employee Markme excluded.Employee Markme, Description excluded.Description, Leave excluded.Leave; try: postgres_hook PostgresHook(postgres_conn_idtutorial_pg_conn) conn postgres_hook.get_conn() cur conn.cursor() cur.execute(query) conn.commit() return 0 except Exception as e: return 1 [create_employees_table, create_employees_temp_table] get_data() merge_data() dag ProcessEmployees()这段代码体现了 Airflow 3 推荐的TaskFlow API编写范式dag装饰器接收调度参数schedule0 0 * * *表示每天零点运行一次Cron 表达式start_date使用pendulum指定 UTC 时区的起始日期catchupFalse关闭回填避免首次部署时把历史周期的任务全部补跑dagrun_timeoutdatetime.timedelta(minutes60)限定单次运行最长 60 分钟普通 Python 函数用task装饰即成为可调度任务依赖关系通过位移运算符表达[create_employees_table, create_employees_temp_table] get_data() merge_data()即两个建表任务并行执行完成后进入get_data最后执行merge_data——清晰的串并行混合 DAG 结构文件末尾的dag ProcessEmployees()必须存在Airflow 解析器通过模块级dag变量发现该 Dag。保存文件后稍等片刻Dag 就会出现在 UI 中。触发并探索你的 Dag打开 Airflow UI在 Dag 列表中找到process_employees。用开关把它切换为on启用调度然后点击播放按钮手动触发一次运行。你可以在Grid网格视图中观察每个 task 的运行状态点击任意 task 可以查看该 task instance 的日志运行成功后你就拥有了一条完整可用的数据管道从外部世界拉取数据、载入 Postgres并持续保持数据整洁。排查建议若get_data失败先检查容器内/opt/airflow/dags/files/目录权限与网络连通性若 SQL 执行失败重点确认连接tutorial_pg_conn的 Host 是否为postgresDocker 网络中的服务名以及 Database/Login/Password 是否与容器初始化值一致merge_data返回1表示异常被捕获可结合任务日志中的异常堆栈定位 SQL 问题。下一步继续深化的方向恭喜你已经用 Airflow 的核心模式与工具构建了一条真实的管道。以下是一些值得继续探索的方向替换 SQL 提供方SQLExecuteQueryOperator来自common-sql提供方天然支持 MySQL、SQLite 等其他数据库只需更换连接类型与对应 provider 即可复用同一套 Dag 逻辑模块化重构将 Dag 拆分为 TaskGroup或把重复逻辑抽取为可复用的自定义 Operator提高代码复用率增加告警与通知在数据处理完成后添加告警步骤通过 Email、Slack 等渠道发送通知。想深入了解更多可以继续浏览仓库中的 Airflow 官方文档索引、阅读 如何编写自定义 Operator 的指南或在 SQL 相关提供方文档 中查找更多执行模式。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表