ARTICLE DETAIL

资讯详情

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

Apache Airflow 数据管道集成实战:三步把 dbt 和 Airbyte 接进来

Apache Airflow 数据管道集成实战:三步把 dbt 和 Airbyte 接进来 Apache Airflow 数据管道集成实战三步把 dbt 和 Airbyte 接进来【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow 凌晨一点CRM、订单库和第三方 SaaS 的数据都要赶在早会前进仓。过去这事靠几条 shell 脚本加人肉盯哪条脚本挂了得等打开报表才发现。Apache Airflow 数据管道集成解决的正是这件事Airflow 负责数据管道调度Airbyte 负责从各个源里抽数据dbt 负责转换建模。下面用一个小时把三者接起来。一张图看懂分工三个角色各管一段Airflow是指挥。DAG有向无环图本质是一张任务依赖图节点是任务箭头表示谁先谁后。它决定什么时候跑、按什么顺序跑、失败怎么办。Airbyte负责抽取。把数据从各数据源拉进仓库支持 CDCChange Data Capture变更数据捕获即只同步变过的那部分而不是每次全量重拉。dbt负责转换。数据进仓后用 SQL 写模型做清洗、标准化、聚合。Airflow 集成 dbt、Airbyte 与 Airflow 联动都不需要自己写胶水代码两家都提供了现成的 Airflow 插件包触发方式就是各家的 API。从 0 到 1装 → 连 → 编 → 跑 先装对两个 Provider 包Provider 就是插件包。Airflow 核心默认不带任何第三方工具集成想用 Airbyte 或 dbt得先把对应的 Provider 装上pip install apache-airflow-providers-airbyte apache-airflow-providers-dbt-cloud 连两个连接名字别起错在 Airflow Web UI 的 Admin → Connections 里建两条airbyte_default类型airbyte填自托管 Airbyte 服务器的地址和 API tokendbt_cloud_default类型dbt_cloud填账号 ID 和 API key这两个名字是 Provider 的默认连接名起别的名字就得在代码里显式传参。连接字段细节可看 Airbyte Provider 文档 和 dbt Cloud Provider 文档。✍️ 放一个 DAG 文件跑通第一次把.py文件丢进dags目录用airflow dags list-import-errors检查有没有导入报错为空再手动 Trigger 一次试跑。Airbyte 同步完成、dbt 作业转绿说明链路通了。一条最小可用 DAG先同步再转换from airflow import DAG from airflow.providers.airbyte.operators.airbyte import AirbyteTriggerSyncOperator from airflow.providers.dbt.cloud.operators.dbt import DbtCloudRunJobOperator with DAG(nightly_pipeline, schedule0 1 * * *, catchupFalse): sync_crm AirbyteTriggerSyncOperator( task_idsync_crm, connection_idAirbyte 连接的 UUID, ) run_dbt DbtCloudRunJobOperator( task_idrun_dbt, job_id12345, ) sync_crm run_dbt依赖关系就一行sync_crm run_dbt转换必须等抽取完成不然 dbt 处理的是半份数据。两个 Operator 默认都是阻塞等待模式——AirbyteTriggerSyncOperator会一直等到同步结束DbtCloudRunJobOperator会等作业跑完才返回最省心得写法。作业特别长时可以换异步模式Operator 提交后立刻把作业 ID 写进 XComAirflow 的任务间传话机制一个任务写、另一个任务取再用 Sensor 任务轮询状态避免一个任务占着执行器空等。 踩坑手记最容易掉的三个坑第一个任务绿了数据没到。现象是 Airflow 显示成功目标表行数却是 0。原因通常是 Airbyte 走增量同步这次根本没捞到新数据源数据没变化或增量字段配错也可能是 dbt 作业里对应模型被跳过了。排查路径别靠猜先看 Airbyte 作业日志里的记录数再看 dbt Cloud 运行结果里实际执行了哪些模型。第二个404Job not found。任务失败报作业不存在。要么job_id填错要么 API key 绑的是另一个 dbt Cloud 账号。对照 dbt Cloud 里该作业页面的 URL 确认 job_id再确认 key 所属账号和作业所在账号是同一个两处都对就不会 404。第三个首跑卡死。第一次同步是全量大表跑几个小时很正常UI 看着像挂起。给任务设上execution_timeout超时就稳妥了——超时直接失败并触发告警好过无限等下去。适用边界哪些管道适合这套组合适合数据源适合定时轮询、批量装载CRM、SaaS、业务库都行转换以 SQL 建模为主团队已经在用 dbt Cloud调度粒度在分钟级以上比如每天凌晨跑一轮不适合、或者要先想清楚需要秒级近实时。Airflow 是批调度它触发的 Airbyte 作业也是批任务真流式要用 Airbyte 自己的持续同步模式自托管 dbt-core 而不是 dbt Cloud。这个 Provider 走的是 Cloud 的 API自托管得换别的方式接没有数仓、转换只是一层简单改名那 dbt 可以先不上省一层复杂度最后一句话记住分工Airflow 指挥、Airbyte 抽数、dbt 建模。下一步就一个动作——挑一个业务上最重要的数据源按上面的四步把它端到端接起来早会报表全绿了再谈扩展。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表