本节目标:把「抽取-转换-加载」从一次性脚本升级成可重跑、有依赖、能落盘分区的数据管道——用标准库手写一个最小 DAG 调度器并真跑,讲透幂等设计,再用 Polars 原生 Parquet 完成分区落盘。
适用版本:Python 3.12+(实测 3.14.6);pandas 3.0.6、polars 2.0.0
11.3 数据管道与 ETL 编排
前两节解决了「单步怎么算得快」,但真实的数据任务从来不是单个脚本——它是一串有先后依赖的步骤:先抽取原始数据,再转换聚合,最后落到数仓或数据集市。而且它要每天/每小时重跑,要能中途失败重试,要能补历史数据。这一节把这些工程问题逐个解决,全部用标准库 + 已装库落地。
11.3.1 ETL 的三段式与它的真正难点
ETL = Extract(抽取)+ Transform(转换)+ Load(加载)。教科书式的描述简单到近乎无用,真正的难点在三个词:
- 可重跑:任务失败重跑一次,结果不能翻倍、不能污染。这叫幂等。
- 有依赖:
transform必须在extract之后,load必须在transform之后。这就是 DAG(有向无环图)。 - 可增量:只处理「新来的」数据,而不是每天全量重算。这叫水位线(watermark)。
现代编排工具(Airflow、Dagster、Prefect)都是在把这三件事做重。理解它们的最好方式是先自己写一个最小的——下面用纯标准库做一个能跑的 DAG 调度器。
11.3.2 幂等的三种手段
幂等(idempotent):同一个任务、同一个输入,跑一次和跑十次,最终结果完全一致。这是数据管道的第一铁律,没有它,重试就是灾难。三种落地手段:
| 手段 | 做法 | 适用 |
|---|---|---|
| 覆盖写 | 每次输出写到同一个确定的路径(按日期命名),后写覆盖先写 | 小批量、可全量重算 |
| 水位线 | 记录「已处理到哪个时间点」,下次只处理其后的数据 | 增量、流式 |
| 去重键 | 落盘时按主键 upsert / deduplicate,重复写入不产生重复行 | 消息类、可能重复投递 |
最容易犯的错是用 append 追加:任务重跑一次,数据就多一份。只要任务是可重跑的,输出就必须是「确定路径 + 覆盖」,或「带主键去重」。 下面调度器里的 load 用的就是覆盖写。
11.3.3 用标准库写一个最小 DAG 调度器
一个 DAG 调度器核心只有两件事:拓扑排序(决定执行顺序)和状态记录(决定哪些跳过)。用 collections.deque 做 Kahn 算法:
import json, time
from pathlib import Path
from dataclasses import dataclass
from collections import defaultdict, deque
WORK = Path("etl"); WORK.mkdir(parents=True, exist_ok=True)
STATE = WORK / "_state.json"
@dataclass
class Task:
name: str
run: callable
deps: tuple = ()
class DAG:
def __init__(self):
self.tasks = {}
def add(self, t: Task):
self.tasks[t.name] = t
def topo_order(self):
indeg = {n: 0 for n in self.tasks}
adj = defaultdict(list)
for n, t in self.tasks.items():
for d in t.deps:
adj[d].append(n); indeg[n] += 1
q = deque([n for n, d in indeg.items() if d == 0])
order = []
while q:
n = q.popleft(); order.append(n)
for m in adj[n]:
indeg[m] -= 1
if indeg[m] == 0:
q.append(m)
if len(order) != len(self.tasks):
raise ValueError("检测到环")
return order
关键点:Kahn 算法结束时若排出的节点数少于总数,说明图里有环——这就是「DAG 校验」。调度器必须在跑之前就拒绝有环的图,而不是跑到一半卡死。
状态记录用一个 JSON 文件(生产里换数据库,机制一样):
def run_dag(dag, run_date):
state = json.loads(STATE.read_text()) if STATE.exists() else {}
done = state.get(run_date, [])
print(f"== 运行 {run_date},已完成: {done}")
for name in dag.topo_order():
if name in done:
print(f" [跳过] {name}(幂等命中)")
continue
t = dag.tasks[name]
t0 = time.perf_counter()
out = t.run(run_date)
print(f" [执行] {name:<18} {(time.perf_counter()-t0)*1000:6.1f} ms -> {out}")
done.append(name)
state[run_date] = done
STATE.write_text(json.dumps(state, ensure_ascii=False))
注意 done 是按 run_date 分桶的:2026-09-30 跑过的任务,不会影响 2026-10-01 的重跑。幂等状态必须带「批次/日期」维度,否则第二天的任务会被当成「已完成」而全部跳过。
三个任务的实现(用 pandas 做转换,落盘为 CSV):
def extract(run_date):
p = WORK / f"raw_{run_date}.csv"
p.write_text("id,city,amount\n1,北京,100\n2,上海,200\n3,北京,300\n", encoding="utf-8")
return p.name
def transform(run_date):
raw = pd.read_csv(WORK / f"raw_{run_date}.csv")
agg = raw.groupby("city", as_index=False)["amount"].sum()
agg.to_csv(WORK / f"agg_{run_date}.csv", index=False) # 确定路径 -> 覆盖写 -> 幂等
return f"{len(agg)} 个城市"
def load(run_date):
agg = pd.read_csv(WORK / f"agg_{run_date}.csv")
agg.to_csv(WORK / f"final_{run_date}.csv", index=False)
return f"落盘 {len(agg)} 行"
dag = DAG()
dag.add(Task("extract", extract))
dag.add(Task("transform", transform, deps=("extract",)))
dag.add(Task("load", load, deps=("transform",)))
真跑(实测输出):
拓扑序: ['extract', 'transform', 'load']
== 运行 2026-09-30,已完成: []
[执行] extract 1.4 ms -> raw_2026-09-30.csv
[执行] transform 37.0 ms -> 2 个城市
[执行] load 4.7 ms -> 落盘 2 行
--- 第二次运行(应全跳过)---
== 运行 2026-09-30,已完成: ['extract', 'transform', 'load']
[跳过] extract(幂等命中)
[跳过] transform(幂等命中)
[跳过] load(幂等命中)
第一次三个任务依次执行,第二次全部跳过——幂等生效。transform 的 37 ms 主要是 pandas read_csv 的固定开销(数据只有 3 行),真实数据量下这部分会被摊薄。
11.3.4 落盘:Parquet 与分区
CSV 落盘是上一节的教训——体积大、要解析、无 schema。生产 ETL 的落盘格式首选 Parquet。50 万行实测(Polars 原生实现):
CSV 大小: 13841 KB Parquet 大小: 3904 KB
写 parquet: 24.0 ms 读 parquet: 31.0 ms
体积约 1/3.5,且列式读取能只取需要的列。 更进一步是分区落盘——按某个低基数列切成多个子目录:
df.write_parquet("partitions/", partition_by="city")
实测生成 8 个 hive 风格目录(city=上海/00000000.parquet、city=北京/... 等)。下游按城市查询时,可以只扫命中的分区:
pl.scan_parquet("partitions/**/*.parquet").filter(pl.col("city") == "北京")
分区键的选择:选低基数、且常作为过滤条件的列(日期、地区、状态)。选错了反而制造「小文件地狱」——几万个几 KB 的碎片文件,比一个大文件还慢。经验值:单分区文件保持在几十 MB 到几百 MB 量级,太小要合并,太大失去裁剪意义。
⚠️ 未实测说明:pandas 的 df.to_parquet() 本机未跑通(无 pyarrow/fastparquet,实测报 ImportError),本节 Parquet 全部用 Polars 原生实现实测;pandas 侧仅保留 CSV 落盘。
11.3.5 从最小调度器到生产编排器
上面这个 60 行的调度器,已经具备生产编排器的两个核心:依赖解析和幂等状态。生产工具在此基础上加的是:
| 能力 | 最小调度器 | Airflow / Dagster |
|---|---|---|
| 依赖 DAG | 拓扑排序 | 同(更复杂的触发规则) |
| 幂等状态 | JSON 文件 | 元数据库(PostgreSQL) |
| 重试 | 无(失败即停) | 指数退避重试 + 超时 |
| 回填 | 无 | backfill 按历史区间批量重跑 |
| 可观测 | print | Web UI + 日志 + 告警 |
| 调度 | 手动触发 | cron / 事件触发 |
注意:Airflow / Dagster / Prefect 本机未安装,上表仅为能力对照,未实测。 但它们的核心机制就是上面这些——理解了你手写的调度器,读它们的文档会快得多。
重试的正确姿势:重试前必须确认任务是幂等的,否则重试会放大错误。tenacity(本机 9.2.1)可以做退避重试,但只包住「可安全重跑」的步骤:
from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, max=10))
def fetch_remote(url):
... # 只重试幂等的 GET;非幂等的写操作要么加去重键,要么别重试
11.3.6 数据质量与可观测
管道跑完「没报错」不等于「数据对」。必须在加载前后加校验,否则错误数据会静默流向下游。最小可用的三条:
def check(run_date):
raw = pd.read_csv(WORK / f"raw_{run_date}.csv")
agg = pd.read_csv(WORK / f"agg_{run_date}.csv")
assert len(raw) > 0, "抽取为空"
assert agg["amount"].sum() == raw["amount"].sum(), "聚合前后金额对不上" # 守恒校验
assert agg["amount"].min() >= 0, "出现负金额"
return "校验通过"
三类校验值得固化成管道步骤:
- 行数/量级校验:本次行数相对上次波动超过阈值(如 ±30%)就告警——这是抓「上游少发了数据」最有效的信号。
- 守恒校验:聚合前后的关键指标(金额、条数)必须相等,能抓住 join 放大、过滤条件写错。
- schema 校验:列名、dtype 变了就报错,别让下游拿到
str却当数字用。
再配上水位线监控:记录每个表「已处理到的时间点」,若某天没推进,说明上游断流。把这三类校验 + 水位线做成管道里的固定节点,数据管道才算真正「可信」。
延伸阅读
- Python 数据工程与 ETL —— 抽取、调度、数仓分层的完整方法论
- Polars 与大规模数据 —— 上一节,Parquet 与分区落盘的底层机制
- 定时任务、幂等与死信处理 —— 第 8 章的幂等与重试,和本节互为表里
小结
- ETL 的真正难点是可重跑(幂等)、有依赖(DAG)、可增量(水位线),不是「读-改-写」本身。
- 幂等三手段:确定路径覆盖写、水位线、主键去重;只要任务可重跑,输出就绝不能是裸
append。 - 一个最小 DAG 调度器只需拓扑排序(Kahn)+ 幂等状态;幂等状态必须带日期/批次维度,否则次日全跳过。
- 实测:首次执行
extract 1.4 ms / transform 37.0 ms / load 4.7 ms,第二次全部跳过。 - 落盘首选 Parquet + 分区:体积约 CSV 的 1/3.5(3904 KB vs 13841 KB),分区键选低基数且常过滤的列,单文件控制在几十到几百 MB。
- pandas 的 Parquet 因本机无 pyarrow 未实测;Airflow/Dagster 未安装,调度器能力对照表为示意。
- 管道必须内建数据质量校验:行数量级、守恒、schema,外加水位线监控——「没报错」不等于「数据对」。
这一章我们从「单机 pandas 怎么不踩坑」,到「Polars 怎么快一个数量级」,再到「把数据任务编排成可信的管道」,走完了数据处理工程的完整链路。下一章转向另一条自动化主线——网络采集:从 HTTP 客户端到页面解析,把外部的、非结构化的数据抓进你的管道。
阅读导航:上一节:Polars 与大规模数据 · 下一节:HTTP 客户端与页面解析 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。