L2.1 数据管道与存储
三维坐标
layer: L2(数据与算子层)|level: Engineer|pillar: 数据模型的上限由数据决定,而数据的可信度由管道决定。本文聚焦「数据进入训练/特征系统之前」的工程化链路:如何用批/流 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 MapReduce → Spark 一统天下,T+1 的离线数仓是标配,「数据延迟一天」被认为天经地义。
- 2016–2021:流处理(Streaming) 崛起,Kafka 成为事实标准的数据总线,Flink 凭借真正的 事件时间(event-time)+ 精确一次(exactly-once) 语义把实时数仓推向主流,「批流一体」成为架构热词。
- 2022 至今:大模型把矛盾推向 「数据质量与可复现性」——训练数据规模爆炸(TB→PB 级),数据契约(Data Contract)、质量门禁(Quality Gate)、数据版本化(DVC/lakeFS) 与 数据血缘(Lineage) 从「锦上添花」变成「不做就翻车」的硬要求。一次脏数据混入预训练语料,代价可能是数十万美元的算力打水漂。
业界信号:
apache/flink、apache/spark、iterative/dvc、great-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 的核心价值之一。