公司动态

Airflow 集成 dbt 与 Airbyte:从零到一搭建每日数据管道的实战指南

📅 2026/8/31 9:03:00
Airflow 集成 dbt 与 Airbyte:从零到一搭建每日数据管道的实战指南
Airflow 集成 dbt 与 Airbyte从零到一搭建每日数据管道的实战指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow这篇文章带你用 Apache Airflow 数据管道串起 Airbyte负责把数据搬进来和 dbt负责把数据洗干净、算明白做一条每日自动跑通的 ETL 链路。适合刚开始接数据任务、还在用手写脚本硬扛的工程师。跟着做一个早上出报表的管道就能落地。手写脚本撑不过三个月痛点出在哪先说结论手写脚本加手动调度撑不过三个月。典型场景是这样的每天早上 7 点老板要一份订单报表。你写了个 Python 脚本凌晨 1 点从业务库抽数据算完发邮箱。第一周很顺。第二周源表改了个字段脚本挂了没人发现。第三周你开始手动盯着执行忘了发日报老板来问。问题不在你在方式。脚本只负责干活没人告诉它几点跑、挂了重不重跑、跑没跑成功。重试只能靠人肉监控等于没有。数据源从 2 个变 5 个之后脚本之间靠记忆衔接一次断点就断一整天。你要的是一个「调度中枢」它按时把每段活派出去记录每一步的结果失败会重试也会喊人。这正是 Airflow 的设计初衷——用代码定义工作流然后交给它去跑。三个工具各管一段Airflow 管什么时候跑Airbyte 管怎么搬dbt 管怎么算这一节把分工讲透。三者各管一段像流水线上的三个工位。Airbyte 负责「搬数据」E 和 L。它是数据集成平台预置了上百种连接器从数据库、SaaS、API 抽数据落到你的数据仓库里。你只需要在它的界面上配好一个 source 到 destination 的同步连接剩下的增量抽取、全量刷新都由它处理。dbt 负责「算数据」T。dbt 是数据转换工具你直接写 SQL 定义模型。raw 层的脏数据经过 staging 层清洗再汇总成 mart 层的业务指标表。它还能对每张表跑测试比如校验主键有没有重复。Airflow 负责「什么时候搬、什么时候算、出了事怎么办」。它不碰数据本身只通过 API 给另外两个下指令到点触发 Airbyte 同步同步完成后再触发 dbt 作业中间失败就重试、就告警。整条链路的数据流长这样Airflow 本身的工作架构你可以参考仓库里的这张官方架构图调度器、worker、元数据库各是什么角色一图看明白安装前先确认版本再装对 Provider代码在 Airflow 里是以「Operator操作器」的形式跑的它们由独立的 Provider 包提供。本例装两个pip install apache-airflow-providers-airbyte pip install apache-airflow-providers-dbt-cloud装之前确认环境Python 3.10这两个 Provider 都要求requires-python 3.10。Provider 包的具体版本写在各自目录下例如providers/airbyte/pyproject.toml声明的是6.0.1providers/dbt/cloud/pyproject.toml声明的是4.9.3装完用pip show对一下即可。在 Airflow 里配两条连接打开 Airflow Web UI进 Connections 新建两条连接Airbyte 连接conn id 建议用默认的airbyte_default。类型选AirbyteHost 填你的服务器地址比如 OSS 部署就是http://localhost:8000/api/v1/Airbyte Cloud 或开了鉴权的部署再补 Client ID 和 Client Secret。字段说明可看 Airbyte 连接文档。dbt Cloud 连接conn id 用dbt_cloud_default。类型选Dbt CloudAPI Token 填你的 User 或 Service Account TokenLogin 字段可以填 Account ID这样后面调 Operator 就不用重复传account_id了。详见 dbt Cloud 连接文档。实战一条每日订单报表管道的四个环节场景固定电商订单每天凌晨从业务库和 SaaS API 同步进数仓dbt 跑一遍模型白天 7 点前 mart 层的日报表必须就绪。下面按「提取 → 转换 → 质检 → 告警」四步拆代码只留关键片段。第一步用 Airbyte 触发同步并等它真正跑完在 Airbyte 界面把订单、用户两张表的同步配好记下那个同步连接的 UUID。DAG 里这样写sync_orders AirbyteTriggerSyncOperator( task_idsync_orders, airbyte_conn_idairbyte_default, connection_id9f2c1a44-7e3d-4b8a-9c6f-1d0e5a2b7c91, # Airbyte 里同步连接的 UUID asynchronousTrue, ) wait_sync AirbyteJobSensor( task_idwait_orders_sync, airbyte_job_idsync_orders.output, timeout3600, poke_interval30, )白话解释AirbyteTriggerSyncOperator只干一件事——按 UUID 触发一次同步。asynchronousTrue表示触发后立刻返回 job id不占着 worker 干等job id 通过.output传给AirbyteJobSensor传感器每 30 秒查一次状态跑完才算这步成功超过 1 小时超时。官方用法和参数见 AirbyteTriggerSyncOperator 文档。第二步同步落地后触发 dbt Cloud 作业同步完成后跑 dbt 模型。job id 在 dbt Cloud 的作业页面查run_dbt DbtCloudRunJobOperator( task_idrun_dbt_models, dbt_cloud_conn_iddbt_cloud_default, job_id30215, check_interval15, timeout1800, )白话解释DbtCloudRunJobOperator调用 dbt Cloud 的 API 发起一次运行然后每 15 秒查一次最长等 30 分钟。这一步结束后staging 和 mart 层的表已经刷新完毕。如果你不想记 job id也可以改用project_nameenvironment_namejob_name三个参数定位作业。更多变体见 dbt Cloud 系统测试示例。第三步质检放在最后用一条 SQL 兜底转换完不能直接交付先验数。最简单的质检是一条行数检查from airflow.providers.standard.operators.python import PythonOperator def check_mart_rows(**ctx): from airflow.providers.standard.hooks.sql import SqlHook hook SqlHook(sql_conn_idwarehouse_default) rows hook.run(SELECT count(*) FROM mart_orders_daily WHERE ds {{ ds }}) if rows 0: raise ValueError(mart 表当天没有数据质检不通过) quality_check PythonOperator( task_idquality_check, python_callablecheck_mart_rows, ) sync_orders wait_sync run_dbt quality_check白话解释查当天分区行数是 0 就直接抛异常DAG 这一步失败后面的交付就不会发生。如果你用 dbt 自带测试更省事——质检那一步换成再触发一次只跑dbt test的 dbt 作业即可写法与第二步相同。第四步失败自动喊人别等早上才发现告警挂在 DAG 级别的失败回调上所有环节挂掉都会触发from airflow.providers.slack.notifications.slack import SlackNotifier def on_pipeline_failed(context): message ( fDAG {context[dag].dag_id} 运行失败 f任务 {context[task_instance].task_id} 需要人工检查 ) SlackNotifier( slack_conn_idslack_default, channel#data-alerts, textmessage, ).notify(context) with DAG( dag_iddaily_orders_pipeline, schedule30 1 * * *, # 每天 01:30 起跑 start_datedatetime(2025, 1, 1), catchupFalse, on_failure_callbackon_pipeline_failed, default_args{retries: 1, retry_delay: timedelta(minutes5)}, ) as dag: ...白话解释on_failure_callback在任务重试耗尽仍失败时触发往 Slack 的#data-alerts频道发一条带 DAG 名和任务名的消息。Slack 连接在 UI 里先建好conn id 为slack_default。到这里凌晨 1:30 起跑、7 点前出数、出事有人收消息的闭环就完整了。三个最容易踩的坑 都是真实会碰到的每条按「现象 → 原因 → 处理」过一遍。坑一把两个 connection id 搞混触发失败。现象Operator 一执行就报连接相关错误。 原因参数名太像。airbyte_conn_id是 Airflow 这边的连接 id如airbyte_defaultconnection_id是 Airbyte 里那个同步连接的 UUID两者缺一不可、各管一头。 处理建连接时把名字起清楚写 DAG 时对着 Airbyte 文档 的字段说明再核一遍。坑二异步模式忘了加 Sensor任务「假成功」。现象DAG 显示全绿但数仓里的表其实没更新完。 原因asynchronousTrue时 Operator 触发完就返回同步还在 Airbyte 那边跑。 处理必须接一个AirbyteJobSensor等真实结果如果就想同步等把asynchronous去掉让 Operator 自己盯状态。坑三重跑触发重复同步或超时设置过短。现象任务偶发超时失败人工重跑后数据量翻倍。 原因触发类任务本身没有幂等保证文档里明确写了不保证 idempotency重试等于再触发一次同时同步时间波动大固定短超时必然误杀。 处理给 Sensor 留足timeout大表同步按历史耗时的 2 倍以上估retries保持 1 次以内重跑前先确认上一次同步的实际状态必要时在 Airbyte 界面手动核对 job 记录。这套方案适合谁下一步学什么这套「Airflow Airbyte dbt」的组合适合每天要出数据、数据源在 2 个以上、又不想再靠人肉盯脚本的小团队。上手前准备好三样一个能访问的 Airflow 环境Python 3.10、Airbyte 实例里配好的同步连接、dbt Cloud 账号和一个可运行的作业。跑通本文这条管道后建议顺着仓库里的 Airbyte 文档目录 和 dbt Cloud 文档目录 把 Sensor、Hook 层的用法补齐再考虑给管道加上更多数据源。先让一条链路稳定跑一周再谈扩展。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考