《Python编程实战》11.3 数据管道与 ETL 编排

把抽取-转换-加载拼成一条可重跑、有依赖的数据管道:讲透幂等的三种落地手段,用标准库手写一个带拓扑排序与幂等状态的 DAG 调度器并真跑,再用 Polars 原生 Parquet 做分区落盘,最后收束成一份数据质量校验、水位线监控与可观测清单。

本节目标:把「抽取-转换-加载」从一次性脚本升级成可重跑、有依赖、能落盘分区的数据管道——用标准库手写一个最小 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 按历史区间批量重跑
可观测printWeb 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 却当数字用。

再配上水位线监控:记录每个表「已处理到的时间点」,若某天没推进,说明上游断流。把这三类校验 + 水位线做成管道里的固定节点,数据管道才算真正「可信」。

延伸阅读

小结

  • 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 客户端与页面解析 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「python」更多文章

  1. 《Python高级编程》目录
  2. 《Python高级编程》11.3 PEP 流程与版本迁移策略
  3. 《Python高级编程》11.2 嵌入式与自由线程运行时