7.2 数据工程
概述与动机
数据工程是构建和运营各种流水线与平台的学科,这些流水线和平台把数据从产生的地方搬运到能创造价值的地方。它涵盖了从源系统的采集、转换为干净且经过建模的形式、以具有成本效益的格式存储、对整个流程的编排,以及让这一切都值得信赖的可靠性实践。如果说数据战略决定了应该存在哪些数据、由谁拥有,那么数据工程就是让数据真正流动起来的管道和机器。
对于大型团队而言,这门学科是基础性的。分析、商业智能、产品实验、机器学习和监管报告,全都处于数据流水线的下游。当这些流水线脆弱、缓慢或不透明时,每一个下游职能都会受损。仪表盘显示的是过时的数字。模型是在被污染的特征上训练出来的。审计人员无法还原一个数字是如何产生的。在企业和政府规模上,流水线要处理来自众多源系统的数十亿条记录,一次悄无声息的故障就可能把错误的数据带入决策、支付或公共统计数据之中。
这个领域已经从定制脚本和单体式 ETL(提取、转换、加载)工具,成长为现代数据技术栈:由开放格式连接起来的、模块化的、大量以 SQL 驱动的采集、转换、编排和存储组件。这种模块化既是一份礼物,也是一个陷阱。它让你能够组合各领域最好的工具,但如果没有工程纪律,就会催生出一片没有文档、没有测试、杂乱蔓延的作业。本章介绍的是那些能让流水线在规模化场景下保持幂等、可测试、可观测且成本可控的实践。
关键原则
- 流水线就是软件,理应享有版本控制、测试、评审和 CI/CD。
- 优先选择幂等、可复现、能够安全重新运行的转换逻辑。
- 让数据流可观测:新鲜度、数据量、模式和质量都要被监控。
- 有意识地为消费者对数据建模,而不是直接倾倒原始表。
- 根据真实的延迟需求来选择批处理还是流处理,而不是图新鲜。
- 把存储格式、分区和计算成本当作头等关切来优化。
- 把采集、转换和服务环节分离开来,使各自能够独立演进。
- 要大声地、尽早地失败;一条中断的流水线,比悄悄输出错误数据要安全得多。
建议
有意识地在 ETL 和 ELT 之间做出选择
ETL 在把数据加载到目的地之前先进行转换。ELT(提取、加载、转换)先加载原始数据,再在一个强大的数据仓库或数据湖仓内部进行转换。现代云平台已经把 ELT 变成了默认选择,因为存储便宜、计算又具有弹性,而保留原始数据能让你在逻辑变化或出现缺陷时重新处理。对于分析类工作负载,优先选择 ELT:先落地不可变的原始数据,再在其之上构建分层的转换。只有在隐私、成本或合同约束要求在数据落地之前就清洗或过滤时,才保留加载前转换的做法。
根据延迟需求来设计批处理和流处理流水线
大多数分析需求用按计划运行的批处理流水线就能很好地满足,这类流水线在推理、测试和回填方面都更简单。只有当业务真正需要低延迟数据时(比如欺诈检测、运营告警、实时个性化)才考虑使用流处理。流处理会带来关于顺序、恰好一次语义、迟到数据和状态管理方面的真实复杂性。当你两者都需要时,考虑那些能统一批处理和流处理逻辑的架构,而不是维护两套彼此分化的代码库。要对自己的延迟需求保持诚实。“实时”常常是一个未经审视的愿望,它会让你的成本翻倍。
用明确的依赖关系来编排
使用一个编排工具,把流水线表达为由任务组成、带有明确依赖关系、重试和调度机制的有向无环图(DAG)。这能让你清楚地看到什么已经运行、什么失败了、什么被阻塞了,还能确定性地回填和重新运行。要把依赖关系建立在数据可用性之上,而不仅仅是时钟时间,这样下游作业才会等待上游数据,而不是凭猜测启动。把编排逻辑纳入版本控制,把 DAG 的改动当作代码改动来对待。
为消费而对数据建模
原始表很少适合分析师直接使用。在需要治理良好、可复用、自助式分析的场景中,应用维度建模,把事实和一致维度组织成星型模式。宽的反规范化表(“单一大表”)在特定查询模式下可能表现更好,对某些消费者来说也更简单,代价是数据冗余和灵活性下降。要对转换进行分层:一个原始暂存层、一个经过清洗和统一的核心层,以及面向消费者的数据集市。这种分离让你可以在一个地方修复逻辑,也让消费者能够依赖稳定的接口。
让流水线保持幂等和可测试
设计转换逻辑,使其重新运行时产生相同的结果,而不是重复或污染数据,例如使用以业务标识符为键的确定性 upsert(更新插入)和分区覆写模式。要在多个层面编写测试:针对转换逻辑的单元测试、模式测试,以及断言唯一性、非空键、引用完整性和可接受取值范围等期望的数据测试。把这些测试运行在 CI 中,这样一个糟糕的改动会在到达生产数据之前就被捕获。
为可观测性和可靠性加装埋点
监控数据健康的四个核心信号:新鲜度(数据是否是最新的)、数据量(行数是否在预期范围内)、模式(结构是否发生了意外变化),以及分布(数值是否出现了异常漂移)。当出现异常时发出告警,并路由给负责的团队。像对待服务一样,为数据事故维护操作手册、值班轮换和无责事后复盘。追踪数据血缘,这样当某处出现故障时,你能立刻看到下游受到的影响。
优化存储和成本
使用 Parquet 等列式开放格式,或支持模式演进、时间旅行和高效更新的开放表格式。按你最常用来过滤的列(通常是日期)对数据进行分区,并通过合并小文件来避免小文件数量激增。用分层存储和生命周期策略把冷热数据分开。按流水线和按查询监控计算开销。失控的成本通常来自全表扫描、缺失分区和无边界的重新处理。要把成本当作一个有明确负责人的指标来对待,而不是月度账单上的一个意外惊喜。
权衡:利弊
| 选择 | 优点 | 缺点 | 最适合场景 |
|---|---|---|---|
| ELT(就地转换) | 保留原始数据,存储便宜,可重新处理 | 存储占用大,需要治理 | 云端分析 |
| ETL(加载前转换) | 控制成本,能及早过滤敏感数据 | 丢失原始数据,更难重新处理 | 受监管或受约束的加载场景 |
| 批处理 | 简单、可测试、易于回填 | 延迟更高 | 大多数分析场景 |
| 流处理 | 低延迟、实时响应 | 复杂、成本高、难以测试 | 欺诈检测、运营告警 |
| 星型模式 | 治理良好、可复用、对自助分析友好 | 前期建模投入大 | 共享的商业智能 |
| 宽表 | 对已知查询速度快、简单 | 数据冗余,灵活性较低 | 窄范围的高性能场景 |
主要的权衡在于简单性与延迟、灵活性之间的取舍。批处理和分层星型模式的 ELT 能给你一个可测试、可回填、易于理解的系统,以经济实惠的方式服务大多数需求。流处理、实时和高度反规范化的设计能换来速度和特定的性能表现,但代价是运维复杂性和测试难度的陡增。只有在具体的业务需求值得为此买单时,才引入这种复杂性,并把简单路径保持为你的默认选择。
与团队讨论的问题
你们是否有意识地选择了 ELT 而不是 ETL,并且保留了不可变的原始数据,以便在逻辑变化或缺陷浮现时能够重新处理? 本章的默认选择是 ELT:以低成本落地原始数据,再构建分层转换,因为保留原始数据能让你在数周之后规则变化或缺陷出现时重新运行一切。删除原始数据会断绝这个选项,这是一个常见且代价惨痛的陷阱。ETL 的相竞争理由在受监管或受约束的加载场景中确实成立,那里的隐私、成本或合同条款要求在数据落地之前就进行过滤或脱敏。带上证据:你们过去需要重新处理历史数据的频率是多少,做不到时代价是什么?对于一条必须能把任何数字追溯到源头的政府或企业流水线而言,不可变的原始记录同样是一项可审计性要求,因此这个答案既塑造了你的存储策略,也塑造了你的法律可辩护性。
数据健康的四个信号中,你们实际监控了哪几个,出问题时会呼叫谁? 本章列出了四个值得关注的信号:新鲜度、数据量、模式和分布。许多团队一个都没有监控,只能靠某位高管盯着一个过时的仪表盘才发现故障,而这是最糟糕的检测方式。在企业和政府规模上,一次悄无声息的故障就可能把错误数据带入支付、报告或公共统计数据之中,因此延迟发现的代价体现在信任和金钱上,而不仅仅是返工。带上你们实际的平均检测耗时,以及目前通常是谁最先发现事故。如果答案是“某个消费者”,你就需要把告警路由给负责的团队,并配上操作手册和无责事后复盘,把数据事故与服务中断完全同等对待。
你们的分析师消费的是经过建模、测试过的数据集市,还是你们直接把原始表甩给他们,还美其名曰“自助式”? 本章说得很直接:原始表很少适合分析师使用,把转换分层为原始暂存层、统一核心层和面向消费者的数据集市,能让你在一处修复逻辑,并给消费者提供稳定的接口。相竞争的诱惑是速度,因为用星型模式或有意设计的宽表建模需要前期投入,跳过它很有诱惑力。但直接倾倒原始数据,会把建模成本反复推给每一位分析师,产生相互矛盾的数字和被浪费的时间。带上一个信号:分析师有多大比例的时间花在重塑原始数据上,有多少个团队重复构建了同样的连接查询。如果这个数字很高,就投资建立一个统一核心层,让消费者依赖经过测试、可复用的接口,而不是各自重新发明。
“实时”在哪里真正物有所值,在哪里只是一个未经审视的愿望、正在悄悄加倍你的运维负担? 本章的默认选择是按计划运行的批处理,它在推理、测试和回填方面更简单,流处理只保留给业务真正需要低延迟的场景,比如欺诈检测或运营告警。相竞争的诱惑是声望,以及利益相关者含糊地要求“实时”数据,这在规划会议上听起来很廉价,到了生产环境却变得昂贵,因为流处理会拖入顺序、恰好一次语义、迟到数据和状态管理,还要再维护一套与批处理逻辑保持同步的代码库。把证据带到讨论现场:对你正在运行或提议的每一条流处理流水线,说出它服务的那个决策,以及那个决策实际能容忍的延迟,用分钟或小时来度量,而不是用形容词。对于一个大型企业或政府平台而言,还要加上每一条实时路径的值班和测试成本,因为一条没人能测试、也没人能全天候值守的流处理流水线,是一个包装成功能的可靠性隐患,而诚实的答案往往会把一个“实时”需求还原成一个能满足同样决策的小时级批处理。
今天你们哪些流水线无法被安全地重新运行,要让每一个转换都变成幂等的,需要付出什么? 本章坚持要求幂等、可复现的转换,使用以业务标识符为键的确定性 upsert 和分区覆写模式,这样重新运行会产生相同的结果,而不是重复或污染数据。相竞争的压力是交付速度,因为一个简单的仅追加作业,比一个专门设计成可重复运行的作业上线更快,而这种走捷径的代价会一直隐藏,直到某次故障迫使你在凌晨两点做一次部分重跑,某人把营收算重复了才暴露出来。带上一份具体的清单:列出那些如果从一个失败点重新运行就会污染数据的作业,并估算其中最严重的一个的影响范围。在企业和政府规模上,一次悄无声息的故障就可能把错误数据带入支付、报告或公共统计数据之中,非幂等处理不仅仅是不便,它还会削弱那种让你在规则变化后能重新处理某个周期、同时仍能把每个数字追溯到源头的可审计性,因此为让重跑变得安全而投入返工,是一个管控问题,而不仅仅是一个整洁与否的问题。
你知道每条流水线运行的成本是多少、由谁对这个数字负责,以及你的云账单中有多少来自全表扫描和缺失分区吗? 本章把存储格式、分区和计算开销当作有明确负责人的头等关切,并警告说失控的成本通常可以追溯到全表扫描、缺失分区和无边界的重新处理。相竞争的考量是,成本方面的工作感觉不如交付新功能那么紧迫,因此常常被一拖再拖,直到月度账单变成一个意外,财务部门开始问起工程团队答不上来的问题。带上证据:按流水线和按查询的开销、来自未分区扫描的成本占比,以及应当被合并的小文件数量。对于一个在众多源系统上处理数十亿条记录的大型组织而言,一份没人负责的云账单会在没有任何一个团队觉得该对此负责的情况下持续增长;而在政府场景中,公共支出必须逐项证明其合理性,因此把计算成本归属给一个有名有姓的负责人、并配上一个被跟踪的指标,能把一笔不透明的开销变成一笔被管理的开销,也常常能发现足以资助下一次平台投资的节省空间。
行业视角
创业公司。 速度胜过架构。把数据采集接入一个托管连接器,构建少量纳入版本控制的转换,并用一个能自行重试和回填的轻量级编排工具来运行它们,而不是手工搭建那些会在深夜悄无声息地崩溃的 cron 作业。从第一次提交开始就让每个模型保持幂等,并为空键和行数添加几个廉价的测试,这样一个糟糕的源数据变更会在 CI 中失败,而不是出现在创始人周一的仪表盘上。不要搭建流处理或定制平台:你最稀缺的资源是工程注意力。
小型企业。 没有专职的数据工程师,因此更倾向于购买一套集成好的技术栈,而不是自己拼装一套。一个托管的 ELT 服务加上一个云数据仓库,能在无需专门平台团队维护的情况下为你提供连接器、调度和存储。把这个选择定位为数据卫生习惯,而不是一个流水线项目:了解哪些源系统在为你的报告提供数据,保留原始数据以便一个错误的数字能被追溯和重新处理,并选择成本可预测的工具,避免一次全表扫描就打爆月度预算。
企业。 问题在于众多团队之间、以及来自众多源系统的数十亿条记录之间的一致性。把 ELT 模式、分层的暂存-核心-集市模型,以及四个数据健康信号标准化,让各团队不再各自重新发明脆弱的流水线。在每个模型上强制执行数据测试和 CI,把计算成本归属给负责的团队,并让数据事故走与服务事故相同的值班、操作手册和无责事后复盘纪律,确保一次悄无声息的故障永远不会在无人察觉的情况下抵达仪表盘。
政府。 采购规则、透明度和公共问责制塑造着流水线的设计。落地不可变的原始记录以满足可审计性,在分层、经过测试的阶段中转换它们,并保留完整的血缘信息,使审计人员能够把任何一个已发布的数字追溯回它的源文件,这往往是一项法律要求。幂等处理让你能在规则变化时安全地重新处理某次申报或某个报告周期,而偏好开放格式和可移植的转换代码,能避免你在一份多年期合同中被锁定在单一供应商身上。
示例
创业公司。 一家十人规模的分析创业公司积累了一堆纠缠不清的 cron 作业,它们会在深夜悄无声息地崩溃,有时工程师手动重跑一个作业还会导致行数被重复计算。团队转向了用于数据采集的托管连接器、用于版本控制模型的转换框架,以及一个能自行重试和回填的轻量级编排工具。他们让每个模型都变成了幂等的,并为空键和行数添加了几个测试,这样一个糟糕的源数据变更现在会在 CI 中失败,而不是出现在创始人周一的仪表盘上。
企业。 一家全球零售商用一套 ELT 技术栈替换了数百个手写的提取脚本。托管连接器负责落地原始源数据,一个转换框架在数据湖仓中构建经过测试、纳入版本控制的模型,一个编排工具通过重试和回填来管理依赖关系。数据测试能在模式漂移到达仪表盘之前就捕获来自源系统的变化。分区列式存储大幅降低了查询成本,同时把新鲜度从按天提升到了按小时。
政府。 一个税务机关通过一条治理良好的流水线采集申报和第三方数据,该流水线以不可变原始记录落地以满足可审计性,再在分层、经过测试的阶段中转换它们。幂等处理让他们能在规则变化时安全地重新处理某个申报周期。完整的血缘信息让审计人员能把任何一个计算出的数字追溯回源文件,这是公共问责制下的一项法律要求。
商业理由:动机、投资回报率与总拥有成本
严谨数据工程的投资回报来自可靠性、速度和成本控制。可靠的流水线意味着决策和报告建立在值得信赖的数据之上,从而避免了错误数字带来的昂贵返工和声誉损害。模块化、经过测试的流水线让团队能更快地交付新的数据产品,使下游每一项分析和机器学习投资的价值不断累积。优化存储和计算能直接降低云账单,一旦分区和查询模式确定下来,降幅往往相当可观。
采用的成本包括平台工具、构建经过测试的模块化流水线所需的工程时间,以及把数据当作软件来对待的纪律。要把这与不采用的成本相权衡:只有作者自己才能理解的脆弱定制作业、被高管发现的悄无声息的数据污染、来自全表扫描的云支出膨胀,以及被数据阻塞等待的分析师。向领导层论证时,把数据工程定位为让分析、商业智能和人工智能变得可信、且负担得起的基础。在这里投入不足,就会限制它之上每一项数据倡议的回报上限。
反模式与陷阱
- 流水线被构建为一次性脚本,没有版本控制、测试或评审。
- 非幂等作业,一旦在故障后重新运行就会重复或污染数据。
- 为了声望而采用流处理,而实际上批处理就能满足延迟需求。
- 把原始表直接甩给分析师,还美其名曰“自助式”。
- 没有可观测性,故障是被下游消费者发现的。
- 忽视分区和文件大小管理,直到云账单爆炸。
- 采集、转换和服务环节相互耦合,导致任何改动都无法安全进行。
- 删除原始数据,使得逻辑变化时无法重新处理。
成熟度模型
- 启动:临时脚本和手动运行,没有测试或监控。故障是被下游消费者发现的,作业无法被安全地重新运行,云成本无人管理、无人归属。
- 发展:一些团队已经采用了编排工具,并把基本转换纳入了版本控制,但在整个组织内做法并不一致。存在零星的测试,幂等性参差不齐,出问题的流水线仍然意味着被动的救火。
- 标准化:带有分层暂存-核心-集市模型的 ELT,经过测试、纳入版本控制,是跨团队应用的文档化标准。带有重试和回填的编排依赖关系、在 CI 中运行的数据测试,以及关于星型模式建模和分区的共享约定,在全组织范围内被强制执行,而不是任由各团队自行其是。
- 管理:平台被度量和控制。新鲜度、数据量、模式和分布被监控,告警路由给负责的团队,流水线服务级别协议、平均检测耗时、数据质量通过率,以及按流水线和按查询的计算成本,都被对照基线加以跟踪。回滚和终止阈值依据证据被强制执行,成本和可靠性都有明确的负责人并被对照目标问责。
- 编排:流水线被完全当作软件来对待,拥有 CI/CD、数据契约和能够在消费者之前捕捉到漂移的自动化异常检测。在延迟真正值得为此买单的地方,批处理和流处理逻辑被统一起来,平台持续改进并实现自助化,容量、存储层级和成本会随着工作负载的变化而自适应地重新平衡,使新的数据产品能够在一个稳定的基础上快速上线。
讨论话题
- 在你的技术栈中,“实时”在哪里真正物有所值,在哪里只是一厢情愿?
- 今天你有哪些流水线无法被安全地重新运行,要修复这一点需要做什么?
- 你的云端数据账单中,有多少来自全表扫描和缺失分区?
- 你的分析师消费的是经过建模的数据集市,还是原始表,这让他们付出了什么代价?
- 你检测一次数据事故的平均耗时是多久,通常是谁最先发现它?
- 统一批处理和流处理逻辑,会降低你的维护负担,还是会增加风险?
关键要点
- 把流水线完全当作软件来对待:版本控制、测试、评审、CI/CD 和可观测性。
- 优先选择带有分层、经过测试模型的 ELT;保留原始数据以便重新处理。
- 默认选择批处理,只有在延迟真正值得为此买单时才选择流处理。
- 让转换保持幂等,使重新运行是安全的。
- 用星型模式或有意设计的宽表,为消费者对数据建模。
- 监控新鲜度、数据量、模式和分布,把数据事故当作服务中断一样对待。
- 把存储格式、分区和计算成本当作头等关切来优化。
参考文献与延伸阅读
- Joe Reis and Matt Housley, “Fundamentals of Data Engineering.”
- Ralph Kimball and Margy Ross, “The Data Warehouse Toolkit.”
- Martin Kleppmann, “Designing Data-Intensive Applications.”
- Bill Inmon, “Building the Data Warehouse.”
- James Densmore, “Data Pipelines Pocket Reference.”
- Nathan Marz and James Warren, “Big Data” (Lambda architecture).
- Barr Moses and colleagues, “Data Quality Fundamentals” (data observability).