ETL/ELT 编排与批处理:从 DAG 到可重跑
0. 元信息
- 主题路径:
docs/topics/data-engineering-basics/subtopics/etl-and-batch-pipelines/README.md - 父主题:
data-engineering-basics - 适合对象:会写 Python/SQL、跑过 cron 脚本、希望从”散装脚本”过渡到”可治理的批处理流水线”的工程师
- 建议周期:1.5~2 周,每周 10~14 小时
- 前置知识:Python、SQL、Docker;建议先完成父主题第 1 阶段
- 最终目标:能为给定业务需求设计一条可重跑、可回填、可观测的 ETL/ELT DAG,并在 Airflow/Prefect/Dagster 三个编排器中至少精通一个
1. 学习路线
DAG 与 Task 模型 → Operator 与 Hook → 调度与回填 → 幂等与状态管理 → 可观测与告警 → 跨编排器对比
2. 阶段周数分配
| 阶段 | 1.5 周方案 | 2 周方案 | 备注 |
|---|---|---|---|
| 1. DAG 模型 | 0.5 天 | 1 天 | 3 任务最小 DAG |
| 2. Operator 与 Hook | 1 天 | 1.5 天 | 5 任务 DAG(CSV→校验→转换→入仓→通知) |
| 3. 调度与回填 | 1 天 | 1.5 天 | cron + catchup + backfill |
| 4. 幂等与状态 | 1 天 | 1.5 天 | 按日期分区 + delete-then-insert |
| 5. 可观测与告警 | 0.5 天 | 1 天 | SLI/SLO + Slack/PagerDuty |
| 6. 跨编排器对比 | 0.5 天 | 1 天 | Airflow / Dagster / Prefect 互译 |
每天 1.5~2 小时。1.5 周方案专注前 5 阶段;2 周方案多一天做跨编排器对比与故障演练。
3. 核心知识 / 产出 / 标准表
| 阶段 | 核心知识 | 实践产出 | 可观察学会标准 |
|---|---|---|---|
| 1. DAG 模型 | 有向无环图、Task、依赖、拓扑序 | 一条 3 任务 DAG | 能用 Python 画拓扑并解释执行序 |
| 2. Operator 与 Hook | BashOperator / PythonOperator / PostgresOperator / S3Hook | 一条 5 任务 DAG(CSV → 校验 → 转换 → 入仓 → 通知) | 能讲清 Operator 与 Hook 边界 |
| 3. 调度与回填 | cron 表达式、catchup、backfill、logical_date | 每日 02:00 调度的 ETL + 一次回填演练 | 能用 airflow dags backfill 重跑历史 |
| 4. 幂等与状态 | 幂等键、分区覆盖、checkpoint、XCom | 一份幂等 DAG(按日期分区 + delete-then-insert) | 能复现重跑产出相同结果 |
| 5. 可观测 | 日志、Metrics、Slack/邮件告警、SLA miss | 一份 SLA 文档 + 告警通道 | 能解释 SLI / SLO 在 ETL 中的体现 |
| 6. 跨编排器对比 | Airflow / Prefect / Dagster / Argo / Temporal 哲学 | 一份对比表 + 至少两个编排器的”Hello World” | 能解释 task-based vs asset-based |
4. 第一周(每天 1.5~2 小时)
Day 1 约定:本计划以 Linux + Docker + 8G 内存为最小实验环境。所有流水线在 Docker Compose 内启动,代码 SQL 配置进 Git。所有任务必须可重跑:清空状态后再跑产出相同结果。
| 日 | 任务 | 当天交付 | 自检 |
|---|---|---|---|
| Day 1 | 装 Docker Compose + Airflow 2.x;跑通最小 Airflow(1 DAG、2 Task:echo hello → echo done) | docker-compose.yml + 第一张 DAG | docker compose up -d 后访问 localhost:8080 看到 DAG;airflow dags list 输出含 hello_dag |
| Day 2 | 写一条 ETL DAG:CSV → 校验 → 转换 → 入仓(PostgreSQL)→ 通知;用 BashOperator / PythonOperator / PostgresOperator | 5 任务 DAG 文件 | airflow dags test etl_orders 全部 success;SELECT COUNT(*) FROM orders 与 CSV 行数一致 |
| Day 3 | 把 DAG 升级为 ELT:先入仓(Parquet → PostgreSQL),再用 PythonOperator 调 dbt 做转换;dbt 加 5 条测试(unique / not_null / accepted_values / relationships / freshness) | DAG + dbt 项目 | dbt build 全部通过;dbt test 输出 5 passed |
| Day 4 | 调度与回填:DAG 设为每日 02:00 schedule="0 2 * * *";用 airflow dags backfill --start-date 2026-01-01 --end-date 2026-01-07 跑一次回填 | 调度配置 + backfill 日志 | airflow dags backfill 7 天全部 success;表行数与 7 天 CSV 累计一致 |
| Day 5 | 幂等性验证:同一日期跑两次 airflow dags trigger etl_orders -e 2026-01-03;验证行数不变;用 EXPLAIN ANALYZE 看 DELETE + INSERT 的 SQL 计划 | 幂等测试报告 | 两次触发后 SELECT COUNT(*) FROM orders WHERE order_date='2026-01-03' 结果相同;EXPLAIN ANALYZE 显示先 DELETE 再 INSERT |
| Day 6 | 可观测:用 StatsD + Prometheus 暴露 DAG duration metric;配 SLA miss 告警(sla=timedelta(minutes=30));用 on_failure_callback 发 Slack 通知 | SLA 文档 + 告警 channel | 故意让 Task 慢 31 分钟触发 SLA miss;Slack 收到告警 |
| Day 7 | 步骤 A:跑通订单 ETL 最小流水线(CSV → PG ODS → dbt DWD → DWS → ADS);步骤 B:补齐 4 类边界(CSV 缺失列、Schema 漂移、分区重复、dbt 测试失败) | 流水线 + 故障演练报告 | pytest -k test_etl_idempotent 4 个用例全过;每类故障复现 + 防护措施写到 notes/week1-day7.md |
Day 7 执行次序
步骤 A —— 最小可用流水线(60~90 分钟)
- CSV → PG ODS(
PythonOperator+to_sql); - dbt 写
stg_orders.sql(清洗)、dwd_orders.sql(明细)、dws_orders_daily.sql(每日汇总); - 三层 dbt 全部
dbt build通过; - 验证:行数从 ODS 到 ADS 逐层变化合理(
dbt show或psql)。
步骤 B —— 4 类边界用例(30~45 分钟)
| # | 用例 | 期望行为 | 验证命令 |
|---|---|---|---|
| B1 | CSV 缺列(如缺 amount) | DAG 失败 + 告警 | pytest -k test_missing_column |
| B2 | CSV 多一列(Schema 漂移) | DAG 继续 + 警告日志 | pytest -k test_extra_column |
| B3 | 同一天重复触发 | 行数不变(幂等) | pytest -k test_idempotent |
| B4 | dbt 测试失败(如 unique 冲突) | DAG 失败 + 阻塞下游 | pytest -k test_dbt_failure |
Day 7 当天必完成步骤 A;步骤 B 至少完成 B1、B3。
5. 阶段通用验收
- 不看答案独立重写一条 5 任务 ETL DAG + 一次 backfill 演练;
- 用自己的话解释”为什么 DAG 必须是 DAG""为什么需要 idempotency""为什么 backfill 与 catchup 是不同操作”;
- 画一张图:DAG 拓扑序图 + 重试状态机 + 数据流图(三件套之一);
- 测试空 CSV、列缺失、Schema 漂移、分区重复、Task 失败 5 类边界;
- 准备至少 3 组自定义 CSV 数据并贴出实际行数与耗时;
- 记录 DAG 端到端 P50 / P99 耗时、Task 失败率、backfill 时长、告警响应时间 4 个核心指标;
- 能修改已有 DAG(加一层、加一个 Task、换一种 Operator)并用同样指标验证。
交付存放:第 3 项的图、第 5 项的输出、第 6 项的指标,统一存到
week1/notes/或week1/<dag>/README.md。
6. 最终验收
-
独立设计并交付一条可重跑、可回填、可观测的 ETL/ELT DAG,覆盖至少 3 个分层、5 类数据测试、2 种告警通道;
-
至少完成 18 个实战项(分布建议):
子阶段 题目数量 难度 平台建议 DAG 模型与 Operator 3 Easy Astronomer Registry 调度与回填 3 Easy / Medium Airflow Tutorial 幂等与状态管理 4 Medium Astronomer Registry + 自己改写 可观测与告警 4 Medium StatsD + Prometheus + Slack 跨编排器对比 4 Medium / Hard Dagster + Prefect + Airflow 互译 约束:至少 9 项达到 Medium,至少 3 项达到 Hard;每项必须留可运行的代码与说明文档。
-
完成 1 个综合项目:mini 数仓 ETL(CSV → Iceberg ODS → dbt DWD/DWS/ADS),含分区 + 幂等 + 测试 + 告警;
-
能用 15 分钟讲清 DAG 拓扑序、Operator 与 Hook 边界、idempotency 实现、backfill 与 catchup 区别、SLI/SLO 设定。
7. 综合项目
首选:mini 数仓 + ETL + dashboard(必做:分层数仓 + dbt 转换 + Streamlit 看板)。
备选:CDC + 实时 dashboard(Debezium → Kafka → Flink → Postgres → Superset)。
备选:把遗留 cron 脚本改造成 Airflow DAG(含幂等、测试、告警)。
mini 数仓必做要求:
- 输入:订单 CSV(用户 / 商品 / 金额 / 时间),按天分区,至少 3 份不同 schema 的 CSV;
- 输出:必输出(1)三层 schema:ODS(原始)→ DWD(清洗 + SCD Type 2)→ DWS(每日汇总)→ ADS(Top 商品 / 用户 ARPU);(2)Streamlit 看板:日订单量 / GMV / Top 10 商品;(3)dbt tests 5 条(unique / not_null / relationships / freshness / accepted_values);(4)Airflow 告警通道(Slack 或邮件);
- 算法 / 工程:DAG 至少 7 Task(extract → validate → load_ods → dbt_run → dbt_test → load_ads → notify);用
EXPLAIN ANALYZE看每条关键 SQL 计划; - 进阶可选:接 Great Expectations 加 5 条分布异常检测;接 OpenLineage 暴露血缘事件。
任何综合项目都必须包含:
- 需求说明与数据契约(源 schema、目标 schema、SLA);
- DAG 架构图与命名规范;
- 核心代码(Airflow DAG + dbt models + 告警 callback);
- 边界测试(源缺失、Schema 漂移、分区重复、Task 失败、dbt 测试失败);
- 可观测(StatsD 指标 + Prometheus 抓取 + Grafana 看板);
- README(设计取舍、已知限制、下一步);
- 复盘记录(模拟 1 次失败,写
retrospective.md); - notes/ 规范:
| 文件 | 内容 |
|---|---|
notes/design.md | DAG 拓扑图、Operator 选型理由、命名规范 |
notes/test.md | 每组测试的输入 / 期望 / 实际 / 通过情况 |
notes/retrospective.md | 用时、难点、收获、改进点 |
notes/runbook.md | 常见失败 + 处置 SOP(Task 失败 / 幂等破环 / 告警风暴) |
notes/dag-screenshots/ | Airflow UI 截图(DAG 树 / Gantt / Task 耗时) |
notes/sql-explain/ | 关键 SQL 的 EXPLAIN ANALYZE 输出 |
notes/data-contracts/ | 数据契约 YAML / Protobuf 定义 |
notes/sample-data/ | 至少 3 组测试 CSV + 行数与耗时表 |
本主题贡献
ETL/ELT 编排的核心矛盾是”批处理天然不幂等 vs 业务要求可重跑可回填”。本子主题专门讲清 Airflow DAG 如何把这条矛盾变成可治理流程——以 Airflow DAG 为编排骨架,用 partition-overwrite 把 Task 变幂等,用 dbt + Spark 做转换,把 SLA 与告警接到 StatsD/Prometheus;不重复父主题的 4 层分层与通用 dbt 介绍。
3 职责
- 用 Airflow DAG TaskFlow API 编排
extract → validate → load_ods → dbt_run → dbt_test → load_ads → notify7 步,把 cron 脚本搬进有向无环图,每个 Task 配retries=3/retry_delay=timedelta(minutes=5)/sla=timedelta(minutes=30)。 - 用 partition-overwrite + 幂等键(按
logical_date分区 delete-then-insert)保证同一日期多次触发产出相同;下游 dbt incremental 与 Spark 批读按日期分区对齐,Task 间不传大数据走 XCom 写 S3。 - 用 StatsD + Prometheus 暴露 DAG duration / task failure rate / SLA miss 三类指标,配
on_failure_callback发 Slack 告警,把批处理流水线接进 on-call 体系。
4 交付物
- 一条 7 任务 Airflow DAG(extract_csv → validate_schema → load_ods → dbt_run → dbt_test → load_ads → notify),TaskFlow API + PostgresOperator + PythonOperator 混用,配
schedule="0 2 * * *"与catchup=False。 - 一份 partition-overwrite 幂等脚本(
DELETE FROM orders WHERE order_date = {{ ds }}后再INSERT),含airflow dags backfill --start-date --end-date演练命令与”重跑 2 次不变行数”验证。 - 一份 dbt 项目(staging / intermediate / marts 三层),5 条 dbt test(unique / not_null / accepted_values / relationships / freshness),
dbt buildexit code 0 作为下游 Task 的成功门槛。 - 一份 SLA 文档 + Slack 告警 callback(含 PagerDuty 路由),4 个核心指标基线:DAG P50/P99 耗时 / Task 失败率 / backfill 时长 / 告警响应时间。
3 指标
- 重跑幂等率 100%(同一
logical_date触发 2 次,COUNT(*)与聚合值完全一致)。 - DAG 端到端 P99 耗时 < 30 min(按数据量级配
executor+pool+priority_weight)。 - SLA miss 告警到达 Slack < 1 min(
on_failure_callback链路耗时,含 webhook 重试)。
8. 推荐开源资料
| 阶段 | 角色 | 资料 | 链接 | 用法 |
|---|---|---|---|---|
| 全部 | 主线书 | Harenslak & de Ruijter《Data Pipelines with Apache Airflow》 | https://www.manning.com/books/data-pipelines-with-apache-airflow | Airflow 实战必读 1~12 章 |
| 全部 | 哲学 | Maxime Beauchemin《Functional Data Engineering》 | https://maximebeauchemin.com/ | 不可变数据 + 函数式流水线 |
| 1~2 | 编排器 | Apache Airflow 官方文档 | https://airflow.apache.org/docs/ | DAG / Operator / Executor |
| 1~2 | 编排器 | Dagster 官方文档 | https://docs.dagster.io/ | Software-defined assets |
| 1~2 | 编排器 | Prefect 官方文档 | https://docs.prefect.io/ | 动态工作流 |
| 1~2 | 模板库 | Astronomer Registry | https://registry.astronomer.io/ | 找 ETL 模板直接复用 |
| 3 | 调度 | crontab.guru | https://crontab.guru/ | cron 表达式验证 |
| 4 | 幂等 | Kleppmann《DDIA》Ch.7 | https://dataintensive.net/ | 事务与幂等 |
| 5 | 可观测 | Prometheus 官方文档 | https://prometheus.io/docs/ | Metrics + Alertmanager |
| 5 | 告警 | StatsD 官方文档 | https://github.com/statsd/statsd | 指标聚合 |
| 6 | 对比 | Lerner《Cloud Native Data Pipelines》 | https://www.oreilly.com/library/view/cloud-native-data/9781098119737/ | 多编排器对比 |
| 6 | 对比 | LakeFS 博客 | https://lakefs.io/blog/ | 数据版本与编排 |
许可证提示:Astronomer Registry 中 DAG 模板以 Apache-2.0 为主;复制或大幅参考前请打开 LICENSE 确认。Airflow、Dagster、Prefect 自身都是 Apache-2.0;商用无需特殊授权。默认做法是读思路后自己重写,而不是复制粘贴。
默认使用顺序:先读 Beauchemin《Functional Data Engineering》建立”不可变数据”心智模型 → 装 Airflow 跑通第一条 DAG → 读 Harenslak 第 1~6 章 → 跑 Astronomer Registry 模板 3 条并改成自己的业务 → 加 dbt 做转换与测试 → 加 StatsD + Prometheus 暴露指标 → 用 crontab.guru 验证 cron → 对比 Dagster / Prefect 的 asset 模型 → 写复盘到 notes/retrospective.md。
9. 学习资料汇聚(v0.3 自包含)
本节由本计划生成。链接指向原始材料或作者公开内容。规范会演进,记录时务必写明版本与日期。
9.1 背景与动机
数据从产生到可用,中间涉及定时拉取、清洗、转换、写入多个步骤。手工 cron 脚本散落各处,重跑靠运气,失败靠人盯。Airflow 2014 年由 Airbnb 开源,把”有向无环图(DAG)“作为批处理核心抽象,引入了依赖、调度、回填、可观测四大支柱。Prefect 2018 年、Dagster 2019 年先后开源,各自提出”动态工作流""Software-defined Assets”等新抽象。今天数据团队在五个编排器中做选择时,必须回答”我要 task-based 还是 asset-based”。
9.2 概念地图
flowchart LR
Cron[cron 脚本] --> DAG[DAG]
DAG --> Task[Task]
Task --> Op[Operator]
Task --> Hook[Hook]
DAG --> Schedule[Schedule cron/catchup]
DAG --> Backfill[Backfill]
Task --> XCom[XCom]
Task --> Idem[幂等键]
DAG --> SLA[SLA miss]
DAG --> Alert[告警]
DAG --> Asset[Asset asset-based]
9.3 基础知识讲解
9.3.1 经典论文与规范
| 资料 | 贡献 | 读法 |
|---|---|---|
| Maxime Beauchemin, Airflow: a workflow management platform(2015) | DAG 作为编排核心抽象 | 看”为什么 DAG” |
| Maxime Beauchemin, Functional Data Engineering(2018) | 不可变数据 + 函数式流水线 | 必读 |
9.3.2 经典书籍
| 书 | 侧重 | 用法 |
|---|---|---|
| Bas P. Harenslak & Julian Rutger de Ruijter, Data Pipelines with Apache Airflow(Manning 2021) | Airflow 实战 | 主线读完 1~12 章 |
| Reuven Lerner, Cloud Native Data Pipelines(O’Reilly 2024) | 多编排器对比 | 选读编排器章节 |
9.3.3 优秀博客与文档
| 资料 | 特点 | 用法 |
|---|---|---|
| Apache Airflow 官方文档 | DAG、Operator、Executor | 查 API 与最佳实践 |
| Astronomer Registry | 官方 DAG 模板库 | 找 ETL 模板 |
| Dagster 官方文档 | Software-defined assets | 看 asset-based 哲学 |
| Prefect 官方文档 | 动态工作流 | 看 flow vs task |
9.3.4 核心人物
| 人物 | 影响 | 材料 |
|---|---|---|
| Maxime Beauchemin | Airflow 起源 | Functional Data Engineering |
| Bolke de Bruijn | Airflow PMC | Apache 邮件列表 |
| Nick Schrock | Dagster 起源 | Dagster 公开演讲 |
| Jeremiah Lowin | Prefect 起源 | Prefect 公开演讲 |
9.3.5 开发方法
| 方法 | 动作 | 何时用 |
|---|---|---|
| Idempotent-by-default | 同一 Task 多次执行产出相同 | 任何 DAG Task |
| Partition-first | 按日期分区,delete-then-insert | 任何 ODS/DWD Task |
| Asset-first(Dagster) | 先定义资产,再写逻辑 | 数据资产视角强的团队 |
| Functional over stateful | 数据不可变、转换是函数 | 借鉴 FDE 思想 |
| SLA-first | 每个 DAG 都有 SLA 与告警 | 生产流水线 |
9.4 经典问题与经典案例(≥5 道)
| # | 问题 | 重要性 | 最简答案 |
|---|---|---|---|
| 1 | DAG 怎么避免循环依赖 | 编排器基本约束 | 用 topological sort 检查 |
| 2 | Task 失败重试怎么设 | 影响可靠性 | 用 retries=3, retry_=timedelta(minutes=5) |
| 3 | 如何做 Backfill | 历史数据重算 | 用 airflow dags backfill --start-date |
| 4 | 跨 DAG 依赖怎么表达 | DAG 间耦合 | 用 TriggerDagRunOperator 或 Sensors |
| 5 | 幂等性怎么保证 | 重跑不重复 | 按日期分区 + delete-then-insert |
| 6 | XCom 传大数据怎么破 | 性能问题 | 不要传大对象,写到 S3/对象存储 |
9.5 学习难点
概念难点
| 难点 | 突破路径 |
|---|---|
| DAG 时间语义 | 用 logical_date 与 data_interval 对照实验 |
| Idempotency vs Transaction | 画 Task 状态机 |
| Backfill vs Catchup | 跑一次 backfill 与一次 catchup |
思维难点
| 难点 | 突破路径 |
|---|---|
| 从需求到 DAG | 按”源 → 转换 → 目标”拆 Task |
| 失败定位 | 看 Task log + SLA miss 标记 |
工程难点
| 难点 | 突破路径 |
|---|---|
| Executor 选错 | 看 LocalExecutor / CeleryExecutor / KubernetesExecutor 对比 |
| 大 DAG 性能 | 用 TaskGroup 拆分 |
9.6 技术标准与接口
Entity
| 名称 | 版本 | 组织 | 状态 | 许可证 |
|---|---|---|---|---|
| Apache Airflow | 2.x | Apache | GA | Apache-2.0 |
| Prefect | 2.x / 3.x | Prefect | 活跃 | Apache-2.0 |
| Dagster | 1.x | Dagster | 活跃 | Apache-2.0 |
| Argo Workflows | 3.x | CNCF | GA | Apache-2.0 |
| Temporal | 1.x | Temporal | 活跃 | MIT |
| Astronomer | 商业 | Astronomer | 活跃 | 商业 |
Scope
Airflow/Prefect/Dagster 是编排器,不替代存储与计算;Argo 偏 Kubernetes-native;Temporal 偏微服务编排。
Structure
- Airflow:
DAG → Task → Operator/Hook → Sensor,Task 间用 XCom 传值。 - Dagster:
Asset → Op → Job,Asset 是数据产物。 - Prefect:
Flow → Task,Flow 是 Python 函数。 - Argo:
WorkflowTemplate → Workflow → DAG → Task。
Ecosystem
- DAG 模板:Astronomer Registry、Airflow Providers。
- 告警:Slack、PagerDuty、Email。
- 监控:Prometheus、StatsD、OpenLineage。
Depth Tiers
| 层级 | 能力 | 标准 |
|---|---|---|
| L0 | 知道存在 | 知道 Airflow/Prefect/Dagster 各自定位 |
| L1 | 看得懂示例 | 能读懂 5 任务 DAG |
| L2 | 能正确调用 | 能写可重跑 DAG 并做 backfill |
| L3 | 能解释与排错 | 能定位 Task 失败、幂等破环、依赖循环 |
| L4 | 能设计与扩展 | 能设计 DAG 框架 + 自定义 Operator |
本子主题目标:L3。
Source
- Apache Airflow 官方文档:引用快照 2026-07-30。
- Dagster 官方文档:引用快照 2026-07-30。
- Astronomer Registry:引用快照 2026-07-30。
9.7 关键代码
9.7.1 Airflow TaskFlow API DAG(5 任务 ETL)
# 9.7.1 Airflow TaskFlow API 写一条 CSV -> Parquet -> PostgreSQL 的 DAG
from datetime import datetime, timedelta
from airflow.decorators import dag, task
import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
@dag(
dag_id="etl_orders_daily",
start_date=datetime(2026, 1, 1),
schedule="0 2 * * *", # 每日 02:00
catchup=False,
retries=3,
retry_=timedelta(minutes=5),
default_args={"owner": "data-eng"},
)
def etl_orders():
@task
def extract() -> str:
df = pd.read_csv("/data/orders.csv")
path = f"/data/orders_{datetime.utcnow():%Y%m%d}.parquet"
pq.write_table(pa.Table.from_pandas(df), path)
return path
@task
def transform(parquet_path: str) -> pd.DataFrame:
df = pq.read_table(parquet_path).to_pandas()
df["total"] = df["qty"] * df["price"]
return df
@task
def load(df: pd.DataFrame) -> int:
# 幂等:先删当月分区再写
from sqlalchemy import create_engine
engine = create_engine("postgresql://airflow:airflow@postgres:5432/dw")
with engine.begin() as conn:
conn.exec_driver_sql(
f"DELETE FROM orders WHERE order_date = '{df['order_date'].iloc[0]}'"
)
df.to_sql("orders", engine, if_exists="append", index=False)
return len(df)
@task
def notify(row_count: int) -> None:
print(f"loaded {row_count} rows")
p = extract()
t = transform(p)
n = load(t)
notify(n)
etl_orders()
9.7.2 Dagster Software-defined Asset
# 9.7.2 Dagster asset-based ETL(同一份代码,不同哲学)
from dagster import asset, Definitions
import pandas as pd
@asset(group_name="etl")
def orders_raw() -> pd.DataFrame:
return pd.read_csv("/data/orders.csv")
@asset(group_name="etl")
def orders_clean(orders_raw: pd.DataFrame) -> pd.DataFrame:
df = orders_raw.copy()
df["total"] = df["qty"] * df["price"]
return df
@asset(group_name="etl")
def orders_summary(orders_clean: pd.DataFrame) -> pd.DataFrame:
return orders_clean.groupby("region")["total"].sum().reset_index()
defs = Definitions(assets=[orders_raw, orders_clean, orders_summary])
9.7.3 Backfill 与幂等性验证脚本
# 9.7.3 用 Airflow CLI 触发 backfill 并验证幂等性
airflow dags backfill etl_orders_daily \
--start-date 2026-01-01 --end-date 2026-01-07
# 验证:同一日期再跑一次,对比 row_count 不变
airflow dags trigger etl_orders_daily -e 2026-01-01
psql -h postgres -U airflow -d dw -c \
"SELECT order_date, COUNT(*) FROM orders GROUP BY 1 ORDER BY 1;"
10. 常见误区
- 把 cron 脚本直接搬进 Airflow,不做幂等;
- Task 用全局变量传值,不走 XCom;
- Backfill 与 catchup 混为一谈;
- Sensor 写死轮询,不配
mode="reschedule"; - SLA 设了不告警;
- Task 失败只重试 1 次就放弃;
- DAG 文件大而全,缺 TaskGroup;
- 用
BashOperator调用 Python 脚本(应直接PythonOperator/@task); - 不挂 OpenLineage,缺血缘;
- 跨编排器随意切换,团队没人精通;
- 误把
data_interval_start/end当execution_date; - Backfill 用
catchup=True而非显式命令。
11. 所有知识点分类
- 编程语言
- 数据结构与算法
- 计算机基础
- 工程技术
- Web 与后端
- 前端与客户端
- 数据与人工智能
- 项目与职业能力
- 安全与可靠性
本计划归属:数据与人工智能 主 + 工程技术 辅。