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.1,providers/dbt/cloud/pyproject.toml声明的是4.9.3,装完用pip show对一下即可。
在 Airflow 里配两条连接
打开 Airflow Web UI,进 Connections 新建两条连接:
Airbyte 连接(conn id 建议用默认的airbyte_default)。类型选Airbyte,Host 填你的服务器地址,比如 OSS 部署就是http://localhost:8000/api/v1/;Airbyte Cloud 或开了鉴权的部署,再补 Client ID 和 Client Secret。字段说明可看 Airbyte 连接文档。
dbt Cloud 连接(conn id 用dbt_cloud_default)。类型选Dbt Cloud,API Token 填你的 User 或 Service Account Token,Login 字段可以填 Account ID,这样后面调 Operator 就不用重复传account_id了。详见 dbt Cloud 连接文档。
实战:一条每日订单报表管道的四个环节
场景固定:电商订单每天凌晨从业务库和 SaaS API 同步进数仓,dbt 跑一遍模型,白天 7 点前 mart 层的日报表必须就绪。下面按「提取 → 转换 → 质检 → 告警」四步拆,代码只留关键片段。
第一步:用 Airbyte 触发同步,并等它真正跑完
在 Airbyte 界面把订单、用户两张表的同步配好,记下那个同步连接的 UUID。DAG 里这样写:
sync_orders = AirbyteTriggerSyncOperator( task_id="sync_orders", airbyte_conn_id="airbyte_default", connection_id="9f2c1a44-7e3d-4b8a-9c6f-1d0e5a2b7c91", # Airbyte 里同步连接的 UUID asynchronous=True, ) wait_sync = AirbyteJobSensor( task_id="wait_orders_sync", airbyte_job_id=sync_orders.output, timeout=3600, poke_interval=30, )白话解释:AirbyteTriggerSyncOperator只干一件事——按 UUID 触发一次同步。asynchronous=True表示触发后立刻返回 job id,不占着 worker 干等;job id 通过.output传给AirbyteJobSensor,传感器每 30 秒查一次状态,跑完才算这步成功,超过 1 小时超时。官方用法和参数见 AirbyteTriggerSyncOperator 文档。
第二步:同步落地后,触发 dbt Cloud 作业
同步完成后跑 dbt 模型。job id 在 dbt Cloud 的作业页面查:
run_dbt = DbtCloudRunJobOperator( task_id="run_dbt_models", dbt_cloud_conn_id="dbt_cloud_default", job_id=30215, check_interval=15, timeout=1800, )白话解释:DbtCloudRunJobOperator调用 dbt Cloud 的 API 发起一次运行,然后每 15 秒查一次,最长等 30 分钟。这一步结束后,staging 和 mart 层的表已经刷新完毕。如果你不想记 job id,也可以改用project_name+environment_name+job_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_id="warehouse_default") rows = hook.run("SELECT count(*) FROM mart_orders_daily WHERE ds = '{{ ds }}'") if rows == 0: raise ValueError("mart 表当天没有数据,质检不通过") quality_check = PythonOperator( task_id="quality_check", python_callable=check_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 = ( f"DAG {context['dag'].dag_id} 运行失败," f"任务 {context['task_instance'].task_id} 需要人工检查" ) SlackNotifier( slack_conn_id="slack_default", channel="#data-alerts", text=message, ).notify(context) with DAG( dag_id="daily_orders_pipeline", schedule="30 1 * * *", # 每天 01:30 起跑 start_date=datetime(2025, 1, 1), catchup=False, on_failure_callback=on_pipeline_failed, default_args={"retries": 1, "retry_delay": timedelta(minutes=5)}, ) 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_default),connection_id是 Airbyte 里那个同步连接的 UUID,两者缺一不可、各管一头。 处理:建连接时把名字起清楚,写 DAG 时对着 Airbyte 文档 的字段说明再核一遍。
坑二:异步模式忘了加 Sensor,任务「假成功」。现象:DAG 显示全绿,但数仓里的表其实没更新完。 原因:asynchronous=True时 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),仅供参考