从零搭建数据管道:用 Airflow、dbt、Airbyte 串起每日数据同步 从零搭建数据管道用 Airflow、dbt、Airbyte 串起每日数据同步【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow每天凌晨 2 点CRM 订单、ERP 库存和支付流水要进数仓。手工跑同步脚本总有人忘第二天早上 BI 看板还是昨天的数据。这套方案用 Apache Airflow 负责调度Airbyte 负责源到仓的数据同步dbt 负责仓内建模把整条管道变成每天自动执行的 DAG。三件套各管一段一句话分工Airflow 管「什么时候跑、失败了怎么办」Airbyte 管「把数据从源搬到仓」dbt 管「把原始表变成能用的模型」。组件职责代表特性适合场景Airflow调度与编排DAG 定义任务依赖重试、超时、池控并发定时批处理、跨系统任务串联Airbyte数据同步内置 500 连接器支持 CDC 增量捕获SaaS 数据、关系型数据库入仓dbt数据建模以 SQL 写模型自带测试与文档数仓分层建模、数据质量校验管道全景整条链路从源系统到消费端如下同步只负责把原始数据完整落到 raw 层不做任何清洗。清洗、打宽、聚合全部交给 dbt 用 SQL 完成模型改动不用碰 Airflow。dbt tests 放在 mart 模型后面作为质量关卡测试不通过任务就失败脏数据不会流到 BI。动手搭建安装 Provider 包Airflow 3.4 起Airbyte 和 dbt Cloud 的能力都在独立 Provider 包里按需安装pip install apache-airflow-providers-airbyte6.0.1 \ apache-airflow-providers-dbt-cloud4.9.3注意本仓库当前 Provider 版本为 airbyte 6.0.1、dbt-cloud 4.9.3Airflow 3.4 要求 Python 3.10。dbt 侧用的是 Cloud 版 Operator如果跑自建 dbt-core装法不同不在本文范围。组件版本基线Python3.10Apache Airflow3.4.0apache-airflow-providers-airbyte6.0.1apache-airflow-providers-dbt-cloud4.9.3配置两个连接Operator 通过 Airflow Connection 拿认证信息用 CLI 加两条连接即可airflow connections add airbyte_default \ --conn-type http \ --conn-host http://airbyte:8000 \ --conn-login api_token --conn-password token airflow connections add dbt_cloud_default \ --conn-type http \ --conn-host https://cloud.getdbt.com \ --conn-password api_tokenairbyte_default和dbt_cloud_default都是 Provider 的默认连接名不填airbyte_conn_id参数时就会用它们。写 Airbyte 同步任务同步任务拆成两个任务Operator 提交同步作业Sensor 轮询直到跑完from datetime import datetime, timedelta from airflow import DAG from airflow.providers.airbyte.operators.airbyte import AirbyteTriggerSyncOperator from airflow.providers.airbyte.sensors.airbyte import AirbyteJobSensor default_args {owner: data-eng, retries: 2, retry_delay: timedelta(minutes5)} with DAG(crm_sync, default_argsdefault_args, schedule0 2 * * *, # 每天凌晨 2 点 start_datedatetime(2026, 1, 1), catchupFalse, tags[sync, airbyte]) as dag: sync_crm AirbyteTriggerSyncOperator( task_idsync_crm, connection_id6e1f..., # Airbyte 里源-目标连接的 UUID asynchronousTrue, # 提交后立即返回 job_id ) wait_crm AirbyteJobSensor( task_idwait_crm, airbyte_job_id{{ ti.xcom_pull(task_idssync_crm) }}, timeout3600, # 最多等 1 小时 poke_interval30, # 每 30 秒查一次 ) sync_crm wait_crm注意connection_id是 Airbyte 侧连接对象的 UUID不是连接名可以在 Airbyte 的 Connection 详情页 URL 里找到。同步量大时建议开deferrableTrue让 Triggerer 代替 Worker 轮询。写 dbt 转换任务模式一样Operator 触发 dbt Cloud 作业Sensor 盯到结束from airflow.providers.dbt.cloud.operators.dbt import DbtCloudRunJobOperator from airflow.providers.dbt.cloud.sensors.dbt import DbtCloudJobRunSensor run_dbt DbtCloudRunJobOperator( task_idrun_dbt, job_id12345, # dbt Cloud 里的 job ID wait_for_terminationFalse, # 不阻塞等待交给 Sensor ) wait_dbt DbtCloudJobRunSensor( task_idwait_dbt, run_id{{ ti.xcom_pull(task_idsrun_dbt) }}, timeout10800, # 最长等 3 小时 poke_interval60, ) run_dbt wait_dbtjob_id也可以不写改用project_nameenvironment_namejob_name三个参数按名字定位好处是 staging 和 prod 环境用同一份 DAG 代码。串成完整 DAG把同步和转换放进同一个 DAG依赖关系一目了然from airflow.operators.empty import EmptyOperator start EmptyOperator(task_idstart) end EmptyOperator(task_idend) start sync_crm wait_crm run_dbt wait_dbt endAirbyte 同步和 dbt 分属两个系统但都在一个 DAG 里表达依赖调度语义比如 Sensor 超时失败会阻止 dbt 运行由 Airflow 统一保证。数据源从 1 个变 5 个时只需在start和run_dbt之间并行加 4 条sync_xx wait_xx链互不阻塞。上生产前的取舍失败回调与重试同步类任务失败多由源端瞬时故障引起重试通常能自愈所以 DAG 级保留retries2。需要人工介入的失败比如 Airbyte 连接配置错了用回调推送出去from datetime import timedelta def notify_failure(context): print(fDAG {context[dag].dag_id} 任务 {context[task_instance].task_id} 失败) # 实际项目里换成企业 IM 或邮件发送 default_args[on_failure_callback] notify_failure default_args[retries] 2 # 失败自动重试 2 次 default_args[retry_delay] timedelta(minutes5) # 间隔 5 分钟回调只做通知和诊断别在里面改数据库状态否则重试机制会放大副作用。变量化配置连接 UUID、dbt job ID 这类值环境间不同硬编码在 DAG 里换环境就要改代码。把它们放进 Airflow Variablesfrom airflow.models import Variable # Web UI 的 Variables 页面写入 # airbyte_crm_connection_id 6e1f... # dbt_mart_job_id 12345 sync_crm AirbyteTriggerSyncOperator( task_idsync_crm, connection_id{{ var.value.airbyte_crm_connection_id }}, # Jinja 模板 )connection_id在 Operator 里是可模板化字段部署时不用改一行代码。池与并发限制Airbyte 和 dbt Cloud 都有并发额度源数据库也扛不住无限并发。用 Pool 给同步任务限流sync_crm AirbyteTriggerSyncOperator( task_idsync_crm, connection_id6e1f..., poolairbyte_pool, # 先在 UI 建一个 2 槽位的 Pool )建议规则airbyte_pool槽位数 源库能承受的并发数一般 2-3dbt_pool槽位数 dbt Cloud 套餐允许的并发作业数。槽位不够时任务排队而不是失败这是池机制的价值。排障速查现象可能原因处理方式同步 Sensor 长时间不结束Airbyte 作业卡在运行态或timeout设得比作业实际耗时短到 Airbyte 后台看作业进度大表同步把timeout提到 7200dbt 任务反复失败上游 raw 表缺数据或 dbt tests 不过先查wait_crm是否真的成功用steps_override[dbt test]单独跑测试定位失败模型多个同步任务互相拖慢源库或 Airbyte 队列并发过高给同步任务加pool限并发错开调度时间到这里一条「凌晨 2 点自动跑、失败会重试、质量不达标会拦下」的管道就搭完了。下一步建议把 dbt 模型测试覆盖率做满再把失败回调接到真实 IM 通道。如果同步量继续涨可以评估 Airbyte 的 CDC 模式和 Airflow 3 的 deferrable 模式让长时间等待的任务不再占用 Worker。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考