跳到主要内容

L2.1 数据管道与存储

三维坐标 layer: L2(数据与算子层)level: Engineerpillar: 数据

模型的上限由数据决定,而数据的可信度由管道决定。本文聚焦「数据进入训练/特征系统之前」的工程化链路:如何用批/流 ETL 把原始数据搬运、清洗、对齐,如何用数据契约质量门禁把脏数据挡在门外,如何对 PII 做脱敏与访问控制,以及如何用 DVC 给数据集打上「可复现的版本号」。这是 L2 数据支柱的入口,下游的特征存储与向量库(L2.2)都建立在这条管道之上。

学习目标

  • 前置知识:读过 L0 全层(尤其 L0.1 全景地图与 L0.2 三大物理墙,知道 AI 系统是「数据→算子→训练→推理」的流水);写过基础 Python 与 SQL,知道「ETL」「数据库」是什么;用过 git。无需 Spark/Flink/Kafka 实战经验,也无需 GPU。
  • 学完产出:① 能画出一条数据从「源→总线→转换→质量门禁→落地」的端到端旅程,并说清每一关的职责;② 能用「有界/无界」一句话讲清批 ETL 与流 ETL 的第一性区别,并按延迟需求为一个真实场景做批/流分工选型;③ 能解释数据契约的兼容性三态(向后/向前/全兼容),判断一次 Schema 变更会不会击穿下游消费者;④ 能说清质量门禁「断言失败即非零退出码阻断管道」的工程铁律,以及 PII 脱敏的可逆/不可逆手段差异;⑤ 亲手用 DVC 给数据集打版本号,并跑通一个「脏数据触发门禁、git+dvc checkout 精确复现历史版本」的最小 lab。
  • 阅读姿势:盯住一条主线——「数据管道的每一道关卡,都是在为下游训练换取可信度:可清洗、可校验、可追溯、可复现」。无论是批流分工、契约校验,还是质量门禁与版本化,本质都在回答同一个问题:当一份数据进了训练集,你能不能证明它「干净且可重放」。

背景与现状

数据管道(Data Pipeline) 的本质,是把散落在业务系统、日志、第三方源里的原始字节,工业化地转化为「结构清晰、质量可保证、版本可追溯」的训练/特征数据。它向下对接千奇百怪的数据源,向上为训练框架(L3)和特征/向量存储(L2.2)提供干净、对齐、可复现的输入——这就是它作为 「数据可信度地基」 的定位。

从产业视角看,AI 数据管道的演进可以概括为三句话:

  • 2015 前:以 批处理(Batch) 为主,Hadoop MapReduceSpark 一统天下,T+1 的离线数仓是标配,「数据延迟一天」被认为天经地义。
  • 2016–2021流处理(Streaming) 崛起,Kafka 成为事实标准的数据总线,Flink 凭借真正的 事件时间(event-time)+ 精确一次(exactly-once) 语义把实时数仓推向主流,「批流一体」成为架构热词。
  • 2022 至今:大模型把矛盾推向 「数据质量与可复现性」——训练数据规模爆炸(TB→PB 级),数据契约(Data Contract)质量门禁(Quality Gate)数据版本化(DVC/lakeFS)数据血缘(Lineage) 从「锦上添花」变成「不做就翻车」的硬要求。一次脏数据混入预训练语料,代价可能是数十万美元的算力打水漂。

业界信号apache/flinkapache/sparkiterative/dvcgreat-expectations/great_expectations 长期高 Star,且 Data Contract 规范(如 Open Data Contract Standard)被多家大厂采纳——这说明「数据可信度」已是大模型团队的胜负手,而非数据工程师的内部事务。

原理与架构

理解 AI 数据管道,最有效的方式是沿着 「一条数据从产生到进入训练集」的旅程,看它穿过了哪些关卡。

2.1 批 ETL vs 流 ETL:两种范式的本质区别

ETL(Extract-Transform-Load,抽取-转换-加载)有两种执行范式,它们不是替代关系,而是按延迟需求分工

维度批 ETL(Spark)流 ETL(Kafka + Flink)
数据边界有界数据集(一个文件/一天分区)无界数据流(永不结束)
延迟分钟~小时级(T+1 / 小时级)毫秒~秒级(近实时)
触发调度器定时触发(Airflow cron)事件到达即触发
典型场景预训练语料清洗、全量特征回填、离线评估集构建实时特征更新、在线日志监控、流式去重
状态管理无状态/读全量有状态(窗口、聚合、join),靠 Checkpoint 容错
一致性语义重跑即幂等exactly-once(Flink 两阶段提交 + Kafka 事务)
代表组件Spark(RDD/DataFrame)Kafka(总线)+ Flink(计算)

核心在于:很多团队的事故源于「拿批的思维写流任务」——例如在流里做无界聚合却不设置窗口,导致状态无限膨胀 OOM。有界/无界是设计的第一性问题

2.2 端到端数据流 + 质量门禁

把抽象的「ETL」落到一条真实的数据流上,质量门禁是其中不可绕过的关卡

这张图的灵魂在 QC(质量门禁):它是一道断路器——数据只有通过「行数下限、关键列非空、数值范围、Schema 匹配」等断言,才允许进入训练集;否则被路由到死信队列(DLQ)隔离并告警。门禁失败时,整条管道应返回非零退出码,让上游 CI/调度器感知失败而不是「静默放行脏数据」。

2.3 数据契约与 Schema 演进

数据契约(Data Contract) 是数据生产方与消费方之间的「API 协议」——它声明字段名、类型、约束、SLA。Schema 演进(Schema Evolution) 则规定契约如何安全变更,核心是 兼容性(Compatibility) 三态:

兼容性类型含义允许的变更升级顺序
向后兼容(Backward)新 Schema 能读旧数据删字段、加带默认值的字段消费者先升级
向前兼容(Forward)旧 Schema 能读新数据加字段、删带默认值的字段生产者先升级
全兼容(Full)同时满足以上两者仅加/删带默认值的字段最安全,推荐

这里的关键在于:在 Kafka 生态里,Schema Registry 强制校验每次 Schema 变更的兼容性,从源头杜绝「生产者改了字段,下游千个消费者集体崩溃」。这本质上是把数据契约做成了编译期检查

2.4 PII 脱敏与访问控制、数据血缘

  • PII 脱敏(Masking / Anonymization):对个人可识别信息(手机号、身份证、邮箱)在进入训练集前做脱敏。常见手段:哈希(不可逆)Tokenization(可逆映射)泛化(年龄→年龄段)差分隐私加噪。原则是「最小可用」——训练用不到的 PII 直接丢弃。
  • 访问控制(Access Control):基于 RBAC/ABAC 对数据列做细粒度授权,敏感列「默认拒绝」,审计每一次访问。
  • 数据血缘(Lineage):记录「这份训练集由哪些源、经哪些转换、在哪个版本产生」。当模型出问题需要溯源,血缘图能在分钟级定位到「是哪批脏数据 / 哪次 Schema 变更」引入的,这正是 DVC 的核心价值之一。

动手实践:用 DVC 做版本化 + 最小质量门禁 ETL

实验目标:在纯 CPU、无需 GPU 的环境下,亲手搭一条「带质量门禁的最小 ETL」并用 DVC 对数据集做版本化——产出物:一个被 DVC 跟踪的数据集 + 一个「质量不达标就返回非零退出码」的 ETL 脚本,建立对「数据可复现 + 质量可阻断」的第一手感知。

路径说明:本实验全程 CPU 可跑,不涉及 GPU/CUDA;只用到 Python 标准库与 dvcpandas

3.1 环境准备

# 推荐 Python 3.11;用 venv 隔离(纯 CPU,无需任何 GPU 依赖)
python3 -m venv .venv && source .venv/bin/activate
pip install "dvc" pandas

# 初始化工作区:git 是 DVC 记录版本的载体
mkdir -p ai-data-lab && cd ai-data-lab
git init -q
dvc init -q
git add .dvc .gitignore && git commit -q -m "chore: dvc init"

3.2 造一份数据集并用 DVC 版本化

# 生成 v1 数据集(含少量脏数据用于演示门禁)
mkdir -p data
cat > data/users.csv <<'CSV'
id,age,country,email
1,28,CN,a@example.com
2,34,US,b@example.com
3,45,CN,c@example.com
4,19,JP,d@example.com
5,52,US,e@example.com
CSV

# 用 DVC 跟踪数据(生成 .dvc 指针文件,真实数据进 DVC 缓存)
dvc add data/users.csv
git add data/users.csv.dvc data/.gitignore
git commit -q -m "data: add users dataset v1"
git tag data-v1

执行 dvc add 后,data/users.csv 的内容进入 .dvc/cache(工作区文件以 link/copy 形式保留、被加入 .gitignore),并生成一个内容寻址的 data/users.csv.dvc 指针文件(含 md5 哈希)。git 只跟踪这个轻量指针,从而实现「代码 + 数据版本一一对应」——这就是数据血缘的基础。

3.3 编写带质量门禁的最小 ETL 脚本

# etl_with_gate.py —— 抽取→脱敏→质量断言;不达标即非零退出(CI 门禁)
import sys
import hashlib
import pandas as pd

SRC = "data/users.csv"
DST = "data/users_clean.csv"

# ---------- Extract ----------
df = pd.read_csv(SRC)

# ---------- Transform:PII 脱敏(email 不可逆哈希)----------
def mask_email(x: str) -> str:
return hashlib.sha256(str(x).encode()).hexdigest()[:12]

df["email"] = df["email"].map(mask_email)

# ---------- Quality Gate:质量断言(行数 / 空值 / 范围 / Schema)----------
def quality_gate(frame: pd.DataFrame) -> list[str]:
errors: list[str] = []
expected_cols = {"id", "age", "country", "email"}

# 1) Schema 检查:列集合必须匹配
if set(frame.columns) != expected_cols:
errors.append(f"schema 不匹配: {set(frame.columns)} != {expected_cols}")

# 2) 行数下限:少于 3 行视为抽取异常
if len(frame) < 3:
errors.append(f"行数不足: {len(frame)} < 3")

# 3) 关键列非空:id / age 不允许空值
for col in ("id", "age"):
n_null = frame[col].isna().sum()
if n_null > 0:
errors.append(f"列 {col} 存在 {n_null} 个空值")

# 4) 数值范围:age 必须落在 [0, 120]
if "age" in frame.columns:
bad = frame[(frame["age"] < 0) | (frame["age"] > 120)]
if len(bad) > 0:
errors.append(f"age 越界行数: {len(bad)}")

return errors


errors = quality_gate(df)
if errors:
print("❌ 质量门禁未通过,阻断管道:", file=sys.stderr)
for e in errors:
print(f" - {e}", file=sys.stderr)
sys.exit(1) # 非零退出码 → CI/调度器感知失败

# ---------- Load ----------
df.to_csv(DST, index=False)
print(f"✅ 质量门禁通过,已写出 {DST}{len(df)} 行,email 已脱敏)")

3.4 运行与观察

# 正常路径:质量通过,写出脱敏后的 clean 数据
python etl_with_gate.py
echo "退出码: $?" # 期望 0

# 把 clean 产物也纳入 DVC 版本化(形成血缘:v1 源 → clean 产物)
dvc add data/users_clean.csv
git add data/users_clean.csv.dvc data/.gitignore
git commit -q -m "data: add cleaned dataset"

# 故障演练:注入脏数据(age 越界 + 删到只剩 2 行)触发门禁
cat > data/users.csv <<'CSV'
id,age,country,email
1,28,CN,a@example.com
2,999,US,b@example.com
CSV
python etl_with_gate.py
echo "退出码: $?" # 期望 1(管道被正确阻断)

# 回滚到干净的 v1 数据集,验证 DVC 可复现性
git checkout data-v1 -- data/users.csv.dvc
dvc checkout data/users.csv # 从缓存还原 v1 真实数据
  • 正常路径:脚本返回退出码 0email 列被替换为 12 位哈希,users_clean.csv 落地。
  • 故障演练:注入「age=999 越界 + 行数不足」后,脚本打印两条门禁错误并返回退出码 1——这就是质量门禁阻断管道的真实效果,上游 CI(如 GitHub Actions)据此判定失败。
  • 可复现性git checkout data-v1 + dvc checkout 能精确还原任意历史版本的数据集,这正是 DVC 把「数据」纳入版本控制的核心收益。

踩坑预警 (Gotchas)

  • dvc add 后别再 git add 原始大文件:DVC 已把真实数据移入缓存,工作区的 .gitignore 会忽略它;误把大文件直接 git add 会让仓库膨胀,违背 DVC 初衷。
  • 质量门禁一定要返回非零退出码:仅 print 警告而不 sys.exit(1),CI/Airflow 会误判任务成功,脏数据照样流入下游——门禁不阻断 = 没有门禁
  • DVC 缓存与远端:本实验数据在本地缓存;生产中需配 dvc remote add(S3/OSS/GCS)并 dvc push,否则换台机器 dvc pull 会失败。
  • PII 脱敏的可逆性陷阱:演示用 SHA256 是不可逆哈希;若业务需「脱敏后还能反查」,应改用带密钥的 Tokenization 而非裸哈希,并严格管控密钥访问权限。

深入思考

下面三题每题先给题干,再用 <details> 折叠一份图文并茂的参考答案。建议先合上答案自己想 3 分钟,再展开对照。

思考题 1:批流分工与 exactly-once

某推荐系统需要「用户点击后 1 秒内更新其兴趣特征」,同时每天凌晨要全量重算用户画像。结合 2.1 节的有界/无界范式区别,你会如何用批/流分工设计这条管道?为什么这两个需求不能用同一套作业解决?流侧如何保证 exactly-once 不重复计费?

展开参考答案(含批流分工架构图 + 算一遍)

结论:1 秒级更新是「无界数据流」问题,只能交给流 ETL(Kafka + Flink);凌晨全量重算是「有界数据集」问题,交给批 ETL(Spark)最划算。两者不是替代而是分工——用同一套作业必然在「延迟」与「吞吐成本」两端都做不好;而 exactly-once 靠 Flink 的 Checkpoint + 两阶段提交 + Kafka 事务三件套,把「读偏移、算状态、写结果」绑成一个原子事务。

为什么不能用同一套作业

  • 流的延迟需求 vs 批的吞吐成本:1 秒级更新要求作业常驻、事件即触发——若用批作业每秒跑一次全量,调度与读全量开销会把成本撑爆;反之全量画像重算要扫一天甚至全历史数据,用流作业逐条算既无必要也无法保证「一次完整快照」。有界/无界是第一性问题(呼应 2.1 节)。
  • 状态语义不同:流侧是有界窗口内的增量状态,批侧是一次性读全量、重跑即幂等。强行合并会落入 2.1 节警示的「拿批思维写流任务」陷阱——无界聚合不设窗口 → 状态无限膨胀 OOM。

exactly-once 怎么保证不重复计费(算一遍依赖链):Flink 周期性打 Checkpoint(如每 30 秒)记录「读到 Kafka 的哪个 offset + 当前算子状态」;写出时用 两阶段提交(2PC) + Kafka 事务——pre-commit 时数据带事务 ID 写入但不可见,待 Checkpoint 全链路成功才 commit 让其可见。一旦作业崩溃,从上个 Checkpoint 恢复:未提交事务回滚、offset 回退重算,同一笔点击的计费结果只会最终可见一次。代价:端到端延迟受 Checkpoint 周期影响(与 2.2 节门禁「失败即阻断」是同一种「宁可慢、不可错」的工程取向)。

思考题 2:Schema 演进事故的兼容性归因

上游团队在 Kafka topic 里给一条消息新增了一个无默认值的必填字段,导致下游所有旧消费者反序列化失败、全线告警。结合 2.3 节的兼容性三态,这违反了哪种兼容性?如果引入 Schema Registry 强制校验,这次事故会在哪一步被拦下?

展开参考答案(含 Schema Registry 拦截时序图 + 对比表)

结论:「加一个无默认值的必填字段」违反了向后兼容(Backward Compatibility)——按新 Schema 升级后的消费者去读旧数据时,旧数据里根本没有这个新字段、又没有默认值可填,反序列化直接失败。引入 Schema Registry(默认 BACKWARD 模式)后,这次变更会在「生产者注册新 Schema」这一步就被拒绝,根本到不了消费者反序列化阶段。

兼容性三态对照(呼应 2.3 节表格):

兼容性谁能读谁「加无默认值必填字段」是否允许
向后兼容 Backward新 Schema 读旧数据❌ 禁止——新消费者读旧数据时无该字段又无默认值可填
向前兼容 Forward旧 Schema 读新数据✅ 允许——旧消费者按旧 Schema 读,忽略未知的新字段
全兼容 Full同时满足两者❌ 仅允许「加/删带默认值的字段」

拦截在哪一步:Schema Registry 把数据契约做成了编译期检查(2.3 节原话)。生产者每次改 Schema 都要先向 Registry 注册并通过兼容性校验;本例因违反 Backward 被直接拒绝注册,生产者拿不到新 Schema ID,脏 Schema 的消息根本写不进 Kafka——事故从「下游 N 个消费者运行时集体崩溃」前移成「上游一次注册被拒、CI 红灯」。正确做法:把该字段改为「带默认值」或分两步演进(先全兼容地加可选字段、等消费者全部升级后再收紧),这与 2.2 节「门禁失败即阻断、不静默放行」是同一治理哲学。

思考题 3:可复现性 vs 存储成本的取舍

团队要求「任意一次模型训练都能精确复现其训练数据」,但数据集达 PB 级,全量多版本存储成本不可承受。结合 DVC 的内容寻址机制与数据血缘,你会如何在「可复现」与「存储成本」之间取舍?哪些数据值得永久版本化,哪些只需记录**生成血缘(怎么算出来的)**即可?

展开参考答案(含可复现性分层决策图 + 算一遍)

结论:可复现不等于「把每个版本的全量数据都存一份」。DVC 的内容寻址天然去重——只有真正变化的内容才占新空间;再叠加「分层策略」:源数据与不可再生数据值得永久版本化,可由确定性流程重算的中间/产物数据只需版本化其「血缘(输入版本 + 代码版本 + 参数)」,用时按需重算,从而用「记录配方」替代「囤积成品」。

算一笔账(数量级估算):假设一份 1 PB 的训练集,每周清洗迭代一次、保留 52 个版本。

  • 朴素全量多版本:52 × 1 PB = 52 PB。按对象存储约 20 美元/TB/月 粗算,52 PB ≈ 52000 TB,每月约 104 万美元——不可承受。
  • DVC 内容寻址去重:若每周仅 ~5% 内容真正变化,实际增量 ≈ 1 PB + 51 × 0.05 PB ≈ 3.55 PB,成本骤降到约 7.1 万美元/月,仅为朴素方案的 ~7%。
  • 再叠加血缘分层:清洗/特征产物(设占 60%)改为「只存血缘 + 按需 dvc repro 重算」,常驻存储进一步压到 ~1.4 PB 量级,月成本降到 3 万美元 上下;代价是复现历史产物要付一次「重算的算力时间」。

取舍原则

  • 值得永久版本化:① 原始采集 / 外部第三方源(已下线就再也拿不回);② 人工标注等高成本不可再生数据;③ 正式发布、用于合规审计的训练集快照。
  • 只记血缘即可:① 由确定性脚本从源数据清洗出的中间产物;② 特征工程、train/val 切分等可重放结果——用 dvc.yaml 的 stage 把「输入数据版本 + 代码 commit + 参数」固化成配方,需要时 dvc repro 重算(呼应 3.2 节「代码+数据版本一一对应」与 2.4 节数据血缘)。核心是把「可复现」的定义从「存下每个成品」松绑为「能确定性重放出每个成品」——这正是 DVC 内容寻址 + 血缘 DAG 的设计意图。

延伸阅读

专题深入:批 ETL 与 DVC 本章已覆盖;训练数据采样/分片/课程学习L2.8Kafka+Flink 流式 ETL 与 Embedding 增量/backfill 管道L2.9

1. 核心 Paper / 规范

  • Apache Flink: Stream and Batch Processing in a Single Engine(2015)— 理解「批流一体」与 event-time/exactly-once 的工程根基。
  • Open Data Contract Standard (ODCS) — 数据契约的开放规范,理解「数据即 API」的协议化思路。
  • Data Management Challenges in Production Machine Learning(Google,2017)— 系统性梳理 ML 数据管道的质量、血缘与版本化挑战。

2. 相关高 Star 仓库与源码必读路径

  • iterative/dvc — 入口看 dvc/repo/add.py(内容寻址)与 dvc/stage/(管道阶段与血缘 DAG)。
  • great-expectations/great_expectations — 看 expectations/ 下的断言体系,对应本文质量门禁的工业级实现。
  • apache/flink — 看 flink-streaming-java 下的 Checkpoint 与两阶段提交,对应 exactly-once 语义。

3. 优质博客 / 视频

  • Confluent 官方博客「Schema Evolution and Compatibility」系列 — 讲透 Kafka 生态的兼容性三态。
  • DVC 官方 Get Started 文档与 lakeFS 博客 — 对比两种数据版本化思路(指针式 vs 类 Git 分支式)。

下一篇L2.2 特征与向量存储:把本文「落地」的数据进一步组织成可低延迟检索的特征与向量,支撑训练与 RAG 在线服务。