公司动态

Apache Airflow 集成实战:如何把 Airbyte 与 dbt 串成一条每天准时跑的数据管道

📅 2026/8/31 19:42:11
Apache Airflow 集成实战:如何把 Airbyte 与 dbt 串成一条每天准时跑的数据管道
Apache Airflow 集成实战如何把 Airbyte 与 dbt 串成一条每天准时跑的数据管道【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 是一个用代码编写、调度和监控工作流的平台。本教程讲 Apache Airflow 集成的实操把 Airbyte 的同步作业和 dbt Cloud 的建模作业挂进同一个 DAG让 T1 日报每天自动跑完、失败可见。从三份日报说起这一节不装依赖先把问题摆清楚为什么需要编排。销售、财务、产品三个团队各要一份 T1 日报销售看订单财务看对账后的收入产品看使用趋势。上游是 CRM、ERP 和事件日志schema 各不一样每天晚上得有人去点同步、去跑建模、去核数。这种人肉编排通常崩在三处有人提前下班下一步没人跑早上九点才发现数据没更新同步明明失败了建模还是照跑产出一批半成品数据定时逻辑散在三个人各自的电脑上改一次时间要喊三遍。你真正想要的是把这一串步骤收进一条管道谁先跑、谁后跑、失败找谁全放在一个能看、能重试、能追溯的地方。这也就是 ETL 工作流编排要解决的问题Airflow 在其中扮演的就是总调度。先分清三件套的分工这一节解决谁干什么避免把三个工具混为一谈或者让 Airflow 重复干别人已经干过的活。一个比较贴切的比喻是工厂Airbyte 是货运把原材料从各家供应商数据源搬到仓库raw 层增量、CDC 这些细节它自己处理dbt 是加工厂把原材料炼成毛坯staging和成品mart顺手做质检测试Airflow 是车间调度不碰货只决定每天凌晨几点开工、哪道工序挂了重跑、每道工序的开始结束时间记在哪里。下面的图把这条流水线画出来箭头方向就是数据的流向上图是 Airflow 自身的组件布局Scheduler 负责什么时候跑Web 界面负责看得见Worker 负责真正执行。本文要记住的只有一件事——Airbyte 的每次同步、dbt 的每次建模最终都会变成 DAG 里的一个任务节点受同一套调度规则约束。动手前装对 Provider、配好两个连接这一节只回答两件事装哪些包、连接里填什么。这两件事没做对后面的代码要么报 ImportError要么一直 401。版本要求与 Airflow 3.x 仓库中各 Provider 已发布版本对齐组件最低版本本文验证版本Python3.93.11Apache Airflow2.103.xapache-airflow-providers-airbyte5.2.36.0.1apache-airflow-providers-dbt-cloud4.4.24.9.3用下面一条命令安装两个 Providerpip install apache-airflow-providers-airbyte5.2.3 apache-airflow-providers-dbt-cloud4.4.2然后在 Web 界面的 Connections 页面建两个连接连接Conn IDHost认证Airbyteairbyte_default按部署形态选填法Client ID/SecretOSS 无认证时留空dbt Clouddbt_cloud_defaulthttps://cloud.getdbt.comAPI Token 填入登录相关字段Airbyte 的 Host 有三种填法最容易填错的就是它部署形态HostAirbyte Cloudhttps://api.airbyte.com/v1/OSS开认证http://localhost:8000/api/public/v1/OSS关认证http://localhost:8000/api/v1/更细的字段说明可以看仓库里的 Airbyte Connection 配置文档包括代理配置和 Token URL。抽取阶段让 Airbyte 同步受 Airflow 调度这一节解决怎么从 DAG 里触发一次同步并且没跑完就算失败。AirbyteTriggerSyncOperator的动作是向 Airbyte 服务器提交一次同步作业拿到job_id然后决定怎么等。等的方式分两种。方式 A同步等待。适合一小时以内能跑完的同步任务自己轮询到终态才结束from datetime import timedelta from airflow.providers.airbyte.operators.airbyte import AirbyteTriggerSyncOperator sync_crm AirbyteTriggerSyncOperator( task_idsync_crm, connection_ida1b2c3d4-xxxx-xxxx-xxxx-xxxxxxxxxxxx, # Airbyte 里源到目的地的连接 UUID airbyte_conn_idairbyte_default, wait_seconds30, # 轮询间隔 timeout5400, # 最多等 1.5 小时 execution_timeouttimedelta(hours1, minutes50), # 超过则取消作业并让任务失败 )这里有个关键区别timeout只限制等多久execution_timeout才是任务的硬上限超过后 provider 会真正调用 Airbyte API 把作业取消掉。两者在 AirbyteTriggerSyncOperator 源码 的execute_complete和on_kill里都能找到对应逻辑。方式 B提交即返回交给 Sensor 盯梢。适合长同步或者一次要并行触发多个同步from airflow.providers.airbyte.sensors.airbyte import AirbyteJobSensor sync_crm AirbyteTriggerSyncOperator( task_idsync_crm, connection_ida1b2c3d4-xxxx-xxxx-xxxx-xxxxxxxxxxxx, asynchronousTrue, # 提交后立刻返回 job_id写入 XCom ) wait_crm AirbyteJobSensor( task_idwait_crm, airbyte_job_id{{ ti.xcom_pull(task_idssync_crm) }}, timeout7200, poke_interval60, ) sync_crm wait_crm如果 worker 数量有限、同步又多又长再给两个算子都加上deferrableTrue等待期间任务会挂起并释放 worker 进程。转换阶段让 Airflow 调度 dbt Cloud 作业这一节解决怎么让 dbt Cloud 里的作业在 DAG 管控下触发、等待、重试。DbtCloudRunJobOperator的逻辑和抽取阶段同构提交一次 job run轮询到终态。from airflow.providers.dbt.cloud.operators.dbt import DbtCloudRunJobOperator run_dbt DbtCloudRunJobOperator( task_idrun_dbt_build, dbt_cloud_conn_iddbt_cloud_default, job_id48617, # dbt Cloud 控制台里的作业 ID check_interval60, # 轮询间隔秒 timeout3600, # 最多等 1 小时 trigger_reasonT1 日报自动触发, execution_timeouttimedelta(hours2), # 超时则取消该 run )几个用得上的点steps_override可以临时覆盖作业里配置好的命令比如夜间只跑[dbt build --select tag:daily]省算力想要提交 等待分离时把wait_for_termination设为False再跟上 sensorfrom airflow.providers.dbt.cloud.sensors.dbt import DbtCloudJobRunSensor wait_dbt DbtCloudJobRunSensor(task_idwait_dbt, run_idrun_dbt.output)retry_from_failureTrue表示如果上一次 run 失败了不新起作业直接按失败 run 的同一份配置重跑省得整条 dbt 流水线从头再来任务在界面里自带 Monitor Job Run 链接点一下就能跳到 dbt Cloud 里对应的 run两边日志对着看很方便。更多参数细节见 dbt Cloud Operators 文档 和 DbtCloudRunJobOperator 源码。拼装阶段接成一条带质量校验的 ETL 工作流这一节回答两阶段怎么拼起来坏数据怎么挡在下游之前。下面是全文唯一的完整示例其余场景直接抄片段即可from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.providers.airbyte.operators.airbyte import AirbyteTriggerSyncOperator from airflow.providers.airbyte.sensors.airbyte import AirbyteJobSensor from airflow.providers.dbt.cloud.operators.dbt import DbtCloudRunJobOperator from airflow.providers.standard.operators.empty import EmptyOperator def check_orders(context): 抽查 mart 层行数异常直接抛错让 DAG 失败 row context[task_instance].conn.execute( SELECT COUNT(*) FROM marts.orders WHERE run_date CURRENT_DATE - 1 ).fetchone() assert row[0] 0, forders mart 无数据: {row[0]} with DAG( dag_iddaily_t1_pipeline, schedule0 2 * * *, start_datedatetime(2026, 1, 1), catchupFalse, default_args{retries: 2, retry_delay: timedelta(minutes10)}, ) as dag: extract AirbyteTriggerSyncOperator( task_idextract, connection_id{{ var.value.crm_conn_id }}, asynchronousTrue, ) wait_extract AirbyteJobSensor( task_idwait_extract, airbyte_job_id{{ ti.xcom_pull(task_idsextract) }}, ) transform DbtCloudRunJobOperator(task_idtransform, job_id48617) quality PythonOperator(task_idquality_check, python_callablecheck_orders) done EmptyOperator(task_iddone) extract wait_extract transform quality done拼装时的三个要点有多个数据源时每个源写一组独立的extract wait_xxx让它们并行再汇合到transform之前connection_id用 Jinja 从 Variables 取换数据源、换环境只改 Variables不动代码质量校验故意卡在transform和done之间断言一挂后面挂 BI 推送之类的任务就不会执行。跑起来之后每个任务节点是成功、在跑还是卡住DAG 视图上一眼可见这正是编排的价值所在。让它跑稳四件必须做的事这一节把经验收进一张表照表取用即可不要过度配置。策略写法收益重试retries2, retry_delaytimedelta(minutes10)可加retry_exponential_backoffTrue吸收偶发网络抖动和限流硬超时 自动取消operator 层execution_timeouttimedelta(hours2)卡死的同步/建模会被 provider 主动 cancel不占坑位并发限流poolairbyte_pool槽位数在 Admin 界面建池时设定一批同步同时触发时不会把 Airbyte 服务器和 worker 打满配置变量化connection_id{{ var.value.crm_conn_id }}换源、换环境只改 VariablesDAG 代码不动重试和超时的分工值得再强调一句retries解决整个任务重跑execution_timeout解决这一轮必须停。只有前者没有后者时一个卡死的同步会反复重试、反复卡死把坑位占满。出事了怎么办四个高频症状的排查路径这一节跳过泛泛的如何排障直接给路径症状 → 先看哪里 → 对应哪个配置项。1. 任务长时间 running但怀疑 Airbyte 侧其实卡住了先开 Airbyte 的 UI 确认作业状态。还在跑说明只是慢调大timeout同时用execution_timeout兜底服务器上根本没在跑则是连接配置问题看下一条。2. 请求一直 401 / 404大概率是 Host 路径或认证填错。三种部署形态的 Host 填法见 Airbyte Connection 配置说明dbt Cloud 侧检查 Token 是否过期、放的位置对不对。3. dbt Cloud 提示已有 run 在进行上一次 run 没到终态新作业堆不上来。给 operator 设reuse_existing_runTrue复用现有 run如果上次是失败用retry_from_failureTrue直接续跑失败的 run。4. 等待型任务占满 worker同步多、每个又长把触发和 sensor 都改成deferrableTrue等待期间任务 defer 出去释放 worker 进程由 trigger 在后台轮询。什么时候该上什么时候别上这套组合适合这样的场景数据是 T1 批处理数据源有两个以上需要统一的调度、重试、超时和监控——此时 Apache Airflow 集成的收益明显大于维护成本。反过来一次性分析脚本直接手动跑一遍就行秒级实时需求不属于这套批处理范式没有仓库、只有几张 CSV 的小团队引入完整编排纯属自找麻烦。最后一句提醒三件套不是捆绑销售。只做转换不做抽取可以只把 dbt 作业挂进 Airflow只做抽取那就只接同步阶段。按你管道里真实存在的搬和加工两步来取用多一步都是多余的。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考