CalcGuide · 技术博客主页 / 一页纸学习计划
🟠

ETL/ELT 编排与批处理:从 DAG 到可重跑

分类:数据与人工智能 · 路径:docs/topics/etl-and-batch-pipelines/README.md

#etl#elt#airflow#prefect#dagster#batch

用 DAG 把数据摄取、转换、写入写成可重跑、可观测、可回填的批处理流水线

父主题

数据工程基础:从 ETL 到数仓与流处理

子主题(0)

ETL/ELT 编排与批处理:从 DAG 到可重跑

0. 元信息

1. 学习路线

DAG 与 Task 模型 → Operator 与 Hook → 调度与回填 → 幂等与状态管理 → 可观测与告警 → 跨编排器对比

2. 阶段周数分配

阶段1.5 周方案2 周方案备注
1. DAG 模型0.5 天1 天3 任务最小 DAG
2. Operator 与 Hook1 天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 与 HookBashOperator / 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 + 第一张 DAGdocker compose up -d 后访问 localhost:8080 看到 DAG;airflow dags list 输出含 hello_dag
Day 2写一条 ETL DAG:CSV → 校验 → 转换 → 入仓(PostgreSQL)→ 通知;用 BashOperator / PythonOperator / PostgresOperator5 任务 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 分钟)

  1. CSV → PG ODS(PythonOperator + to_sql);
  2. dbt 写 stg_orders.sql(清洗)、dwd_orders.sql(明细)、dws_orders_daily.sql(每日汇总);
  3. 三层 dbt 全部 dbt build 通过;
  4. 验证:行数从 ODS 到 ADS 逐层变化合理(dbt showpsql)。

步骤 B —— 4 类边界用例(30~45 分钟)

#用例期望行为验证命令
B1CSV 缺列(如缺 amountDAG 失败 + 告警pytest -k test_missing_column
B2CSV 多一列(Schema 漂移)DAG 继续 + 警告日志pytest -k test_extra_column
B3同一天重复触发行数不变(幂等)pytest -k test_idempotent
B4dbt 测试失败(如 unique 冲突)DAG 失败 + 阻塞下游pytest -k test_dbt_failure

Day 7 当天必完成步骤 A;步骤 B 至少完成 B1、B3。

5. 阶段通用验收

  1. 不看答案独立重写一条 5 任务 ETL DAG + 一次 backfill 演练;
  2. 用自己的话解释”为什么 DAG 必须是 DAG""为什么需要 idempotency""为什么 backfill 与 catchup 是不同操作”;
  3. 画一张图:DAG 拓扑序图 + 重试状态机 + 数据流图(三件套之一);
  4. 测试空 CSV、列缺失、Schema 漂移、分区重复、Task 失败 5 类边界;
  5. 准备至少 3 组自定义 CSV 数据并贴出实际行数与耗时;
  6. 记录 DAG 端到端 P50 / P99 耗时、Task 失败率、backfill 时长、告警响应时间 4 个核心指标;
  7. 能修改已有 DAG(加一层、加一个 Task、换一种 Operator)并用同样指标验证。

交付存放:第 3 项的图、第 5 项的输出、第 6 项的指标,统一存到 week1/notes/week1/<dag>/README.md

6. 最终验收

7. 综合项目

首选:mini 数仓 + ETL + dashboard(必做:分层数仓 + dbt 转换 + Streamlit 看板)。
备选:CDC + 实时 dashboard(Debezium → Kafka → Flink → Postgres → Superset)。
备选:把遗留 cron 脚本改造成 Airflow DAG(含幂等、测试、告警)。

mini 数仓必做要求:

任何综合项目都必须包含:

  1. 需求说明与数据契约(源 schema、目标 schema、SLA);
  2. DAG 架构图与命名规范;
  3. 核心代码(Airflow DAG + dbt models + 告警 callback);
  4. 边界测试(源缺失、Schema 漂移、分区重复、Task 失败、dbt 测试失败);
  5. 可观测(StatsD 指标 + Prometheus 抓取 + Grafana 看板);
  6. README(设计取舍、已知限制、下一步);
  7. 复盘记录(模拟 1 次失败,写 retrospective.md);
  8. notes/ 规范
文件内容
notes/design.mdDAG 拓扑图、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 职责

  1. 用 Airflow DAG TaskFlow API 编排 extract → validate → load_ods → dbt_run → dbt_test → load_ads → notify 7 步,把 cron 脚本搬进有向无环图,每个 Task 配 retries=3 / retry_delay=timedelta(minutes=5) / sla=timedelta(minutes=30)
  2. 用 partition-overwrite + 幂等键(按 logical_date 分区 delete-then-insert)保证同一日期多次触发产出相同;下游 dbt incremental 与 Spark 批读按日期分区对齐,Task 间不传大数据走 XCom 写 S3。
  3. 用 StatsD + Prometheus 暴露 DAG duration / task failure rate / SLA miss 三类指标,配 on_failure_callback 发 Slack 告警,把批处理流水线接进 on-call 体系。

4 交付物

  1. 一条 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
  2. 一份 partition-overwrite 幂等脚本(DELETE FROM orders WHERE order_date = {{ ds }} 后再 INSERT),含 airflow dags backfill --start-date --end-date 演练命令与”重跑 2 次不变行数”验证。
  3. 一份 dbt 项目(staging / intermediate / marts 三层),5 条 dbt test(unique / not_null / accepted_values / relationships / freshness),dbt build exit code 0 作为下游 Task 的成功门槛。
  4. 一份 SLA 文档 + Slack 告警 callback(含 PagerDuty 路由),4 个核心指标基线:DAG P50/P99 耗时 / Task 失败率 / backfill 时长 / 告警响应时间。

3 指标

  1. 重跑幂等率 100%(同一 logical_date 触发 2 次,COUNT(*) 与聚合值完全一致)。
  2. DAG 端到端 P99 耗时 < 30 min(按数据量级配 executor + pool + priority_weight)。
  3. 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-airflowAirflow 实战必读 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 Registryhttps://registry.astronomer.io/找 ETL 模板直接复用
3调度crontab.guruhttps://crontab.guru/cron 表达式验证
4幂等Kleppmann《DDIA》Ch.7https://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 BeaucheminAirflow 起源Functional Data Engineering
Bolke de BruijnAirflow PMCApache 邮件列表
Nick SchrockDagster 起源Dagster 公开演讲
Jeremiah LowinPrefect 起源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 道)

#问题重要性最简答案
1DAG 怎么避免循环依赖编排器基本约束用 topological sort 检查
2Task 失败重试怎么设影响可靠性retries=3, retry_=timedelta(minutes=5)
3如何做 Backfill历史数据重算airflow dags backfill --start-date
4跨 DAG 依赖怎么表达DAG 间耦合TriggerDagRunOperator 或 Sensors
5幂等性怎么保证重跑不重复按日期分区 + delete-then-insert
6XCom 传大数据怎么破性能问题不要传大对象,写到 S3/对象存储

9.5 学习难点

概念难点

难点突破路径
DAG 时间语义logical_datedata_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 Airflow2.xApacheGAApache-2.0
Prefect2.x / 3.xPrefect活跃Apache-2.0
Dagster1.xDagster活跃Apache-2.0
Argo Workflows3.xCNCFGAApache-2.0
Temporal1.xTemporal活跃MIT
Astronomer商业Astronomer活跃商业

Scope

Airflow/Prefect/Dagster 是编排器,不替代存储与计算;Argo 偏 Kubernetes-native;Temporal 偏微服务编排。

Structure

Ecosystem

Depth Tiers

层级能力标准
L0知道存在知道 Airflow/Prefect/Dagster 各自定位
L1看得懂示例能读懂 5 任务 DAG
L2能正确调用能写可重跑 DAG 并做 backfill
L3能解释与排错能定位 Task 失败、幂等破环、依赖循环
L4能设计与扩展能设计 DAG 框架 + 自定义 Operator

本子主题目标:L3

Source

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. 常见误区

11. 所有知识点分类

  1. 编程语言
  2. 数据结构与算法
  3. 计算机基础
  4. 工程技术
  5. Web 与后端
  6. 前端与客户端
  7. 数据与人工智能
  8. 项目与职业能力
  9. 安全与可靠性

本计划归属:数据与人工智能 主 + 工程技术 辅。


直接依赖(2)

查看知识图谱