ddia-part3
《Designing Data-Intensive Applications》(Martin Kleppmann, O’Reilly 2017) 第三部分「派生数据」读书笔记。用我自己的话重新梳理每章的核心概念、关键权衡和例子。第一部分讲单机的地基,第二部分讲多机的分布式难题。第三部分回答一个更高层的问题:现实中没有一个数据库能满足所有需求,我们总要把多个系统(数据库、缓存、搜索索引、数仓、消息队列……)组合起来——那么如何让这些系统之间的数据保持一致、可靠地流动? 这就是”派生数据”的主题。
本篇目录
- 第 10 章 · 批处理 (Batch Processing)
- 第 11 章 · 流处理 (Stream Processing)
- 第 12 章 · 数据系统的未来
先建立一个核心区分——按数据如何流入,系统分三类:
- 在线服务 (services):等着请求来,尽快响应。衡量指标是响应时间(第 1 章讲过)。
- 批处理 (batch):一次性处理大批已积累的数据,跑完产出结果。衡量指标是吞吐量,不在乎单条延迟。
- 流处理 (stream):介于两者之间——数据来一条处理一条,接近实时,但输入是”永不结束”的。
第 10、11 章分别讲批处理和流处理,第 12 章把它们统一起来,畅想整个数据系统生态的未来形态。
第 10 章 · 批处理 (Batch Processing)
从 Unix 哲学讲起
作者从一个看似不相干的地方切入——Unix 命令行。用一行管道就能分析日志:
cat access.log | awk '{print $7}' | sort | uniq -c | sort -rn | head
(取出每行第 7 个字段 → 排序 → 统计去重计数 → 按次数倒序 → 取前几名。)
这条命令体现了 Unix 哲学,而它正是现代批处理系统的思想源头:
- 每个程序只做好一件事 (
sort只管排序,uniq只管去重)。 - 程序之间用统一接口拼接:一切皆为字节流(stdin/stdout),所以任意工具都能用管道
|组合。 - 输入不可变:命令不修改输入文件,只产出新的输出。这意味着你可以随便重跑、随便实验,不怕搞坏原始数据。
💡 思想 tip:记住这三点——单一职责、统一接口、输入不可变。这是整章(甚至整个第三部分)的灵魂。后面的 MapReduce、Spark,本质上都是”把 Unix 管道搬到几千台机器上”。当你理解”为什么不可变输入让系统更健壮”,你就抓住了大数据系统设计的精髓。
MapReduce:分布式版的 Unix 工具
MapReduce(Google 提出,Hadoop 是其开源实现)就是把上面的思路扩展到成千上万台机器。它的输入输出通常存在分布式文件系统 HDFS 上——一种 shared-nothing、每个文件块多副本冗余的存储。
一个 MapReduce 作业要你提供两个纯函数:
- Mapper:对每一条输入记录调用一次,从中提取出 key 和 value。(如:从每行日志提取出
(URL, 1)。) - Reducer:框架会把相同 key 的所有 value 收集到一起,交给 reducer 聚合。(如:把某个 URL 的所有
1加起来 = 访问次数。)
中间最关键、也最耗资源的一步叫 Shuffle(洗牌):框架自动把 mapper 的输出按 key 分区、排序,然后送到对应的 reducer。“把相同 key 的数据搬到一起”就是 shuffle 干的事,它是 MapReduce 的核心机制,也常是性能瓶颈。
Join 怎么做? 批处理里 join 很常见(如”给每条用户活动记录关联上用户的年龄”):
- Reduce 端 join(排序-合并 join):把两份数据都按 join key 排序、shuffle 到一起,在 reducer 里合并。通用但要经过昂贵的 shuffle。
- Map 端 join:如果一份数据小到能装进内存(如用户表),就把它广播到每个 mapper(广播哈希 join),省掉 shuffle,快很多。
MapReduce 的设计哲学:为什么它可靠
- 输入不可变、输出到新位置:作业从不修改输入,失败了直接重跑即可,天然容错、易调试。
- 容错靠重试:某个任务失败,框架在别的机器上重跑它,不影响整体。
- 物化中间结果:每一步的输出都落盘(写进 HDFS),下一步再读。这让它极其健壮,但也慢——大量磁盘 I/O。
超越 MapReduce:数据流引擎
MapReduce “每步都落盘”的笨重催生了新一代数据流引擎 (dataflow engines)——Spark、Flink、Tez。它们把整个作业建模成一张算子的有向无环图 (DAG):
- 不把中间结果落盘,尽量留在内存里、直接流向下一个算子 → 快得多。
- 容错换了思路:不再靠”重读落盘的中间数据”,而是记录数据的血缘 (lineage)——即”这份数据是由哪些输入、经过哪些运算得来的”。某部分丢了,就沿着血缘重新计算。Spark 的 RDD 就是这个思想。
💡 选型 tip:Spark 比 Hadoop MapReduce 快,核心原因不是什么魔法,就是省掉了中间结果反复读写磁盘(用内存 + 血缘重算代替)。理解这一点,你就明白:Spark 在需要多步迭代的任务(如机器学习)上优势巨大,但它对内存要求更高;而经典 MapReduce 在超大规模、内存放不下的场景下反而更稳。又是一次”性能 vs 健壮性/资源”的权衡。
此外还有针对图的批处理模型(如 Pregel 的”整体同步并行 BSP”,让顶点之间反复传消息迭代),以及建在这些引擎之上的高级声明式接口(Hive、Pig、Spark SQL / DataFrame)——你写类 SQL,引擎自动优化成底层的 MapReduce/DAG 作业,呼应了第 2 章”声明式让引擎自由优化”的思想。
批处理的典型产出:构建搜索索引、批量训练机器学习模型、生成推荐、预计算好一批只读的键值数据供在线服务查询。关键原则始终是:把输出写到新地方,绝不原地修改输入。
第 11 章 · 流处理 (Stream Processing)
从”有界”到”无界”
批处理的前提是输入有界(数据集大小固定,跑完就结束)。但现实中数据是持续不断产生的——用户点击、传感器读数、交易……你不想等一天攒够了再批处理,而想数据一到就处理。这就是流处理:处理无界 (unbounded) 的数据。
流里的基本单位是事件 (event):一个小的、自包含的、不可变的记录,记录”某时刻发生了某事”。
一、怎么传输事件流:消息系统
生产者产生事件,消费者处理事件,中间通常靠消息中间件 (message broker)。有两大流派:
- 传统消息队列 (AMQP/JMS 风格,如 RabbitMQ):消息被某个消费者处理确认 (ack) 后就从队列删除。适合”任务分发”——多个消费者分摊负载,每条消息只被处理一次。缺点:消息一旦消费就没了,无法回放。
- 基于日志的消息 (log-based,如 Kafka):消息追加到一个持久的、只增不减的日志里,消费者靠记录自己的偏移量 (offset) 来跟踪读到哪了。
💡 架构 tip:Kafka 这类”基于日志”的设计是流处理的一大关键洞察,它有两个杀手锏。其一,可回放 (replay):因为消息不删除,新加一个消费者可以从头重读所有历史事件——这对”上线一个新的派生视图、需要用全部历史数据重建”至关重要。其二,多消费者独立消费:同一份日志可以被搜索索引、缓存、数仓等多个下游各自按自己的进度读取,互不干扰。这正是把它当作”多个系统之间数据流动的主干”的基础。
二、数据库与流:把数据库变化当成流
这是本章最重要的思想之一。数据库的每一次写入,其实都是一个”事件”。如果我们能把数据库的变更捕获成一条事件流,就能把它同步给所有需要的派生系统。两种做法:
- 变更数据捕获 (CDC, Change Data Capture):读取数据库的复制日志(第 5 章讲过),把每一行的增删改作为事件流出来,用于实时更新搜索索引、缓存、数仓,让它们和主库保持同步。
- 事件溯源 (Event Sourcing):更彻底——不存”当前状态”,而是把所有发生过的事件按顺序、不可变地存下来,当前状态是把这些事件从头”重放”计算出来的结果。
💡 核心洞察 tip:这两者共享一个深刻的思想——日志(事件流)才是”真相之源 (source of truth)”,而数据库里的表、缓存、索引都只是从这个日志派生出来的、可随时重建的”视图”。这个视角非常强大:因为事件不可变、可回放,你随时可以清空某个派生视图、从日志重新构建它(比如换了搜索引擎、缓存结构要改,都不怕);出了 bug 也能追溯每一步是怎么来的。这就是”派生数据”这个部分标题的真正含义——把数据系统重新理解为”一个真相日志 + 一堆从它派生的视图”。
区分两个概念也有用:命令 (command) 是”用户想做某事”的请求(可能被拒绝),一旦验证通过、成为既成事实,就变成不可变的事件 (event)。事件流里存的都是已发生、不可更改的事实。
三、怎么处理流
流处理的用途:复杂事件检测(如”连续三次登录失败就告警”)、实时分析、维护物化视图、实时搜索等。这里有两个难点特别值得讲:
难点一:时间。 流处理里”时间”极其棘手,因为有两个不同的时间:
- 事件时间 (event time):事件实际发生的时间。
- 处理时间 (processing time):事件被系统处理的时间。
两者常常对不上——手机 App 离线时攒下的事件,可能几小时后联网才上报。如果你按”处理时间”统计”每分钟事件数”,一个网络延迟就会让统计全错。更麻烦的是掉队者 (straggler):你以为某分钟的数据都到齐了、算出了结果,结果又迟到来了一条属于那一分钟的事件——怎么办?
处理时间还要用到窗口 (window):把无界的流切成有限的块来聚合。常见有滚动窗口 (tumbling,不重叠)、跳动窗口 (hopping)、滑动窗口 (sliding)、会话窗口 (session)。
💡 实战 tip:“事件时间 vs 处理时间”是流处理里最容易踩的坑,没有之一。 只要你要做任何”按时间聚合”的统计,就必须想清楚:我按哪个时间?迟到的数据怎么办?(要么设一个等待期,要么允许事后修正结果。)很多实时看板数据不准,根子就在这里。
难点二:流上的 join。 有三种:流与流 join(如”点击广告”事件配对”购买”事件)、流与表 join(用事件去关联数据库里的维度信息,需要 CDC 把表也变成流)、表与表 join(维护物化视图)。难点在于两条流的事件到达时间不确定,要维护状态、处理乱序。
容错:exactly-once(恰好一次)。 流是无限的,不能像批处理那样”失败就整个重跑”。做法有微批处理 (microbatching)、检查点 (checkpoint) 等。要真正做到”每个事件的效果恰好生效一次”,靠的是幂等 (idempotence)(重复执行结果不变——正好呼应上一次修 daily 重复 bug 用的手段!)和原子提交。
第 12 章 · 数据系统的未来
最后一章是作者的”登高望远”,把前面所有内容串成一个统一的世界观,并以一段严肃的伦理反思收尾。
一、数据集成:没有银弹,只有组合
现实第一课:没有任何单一工具能满足所有需求。 你几乎必然要同时用关系库(事务)、搜索索引(全文检索)、缓存(加速)、数仓(分析)……问题随之而来——同一份数据存在好几个系统里,怎么保证它们不打架?
作者给的答案,正是第 11 章那个思想的延伸:指定一个”真相之源”(通常是一个事件日志),其他所有系统都作为它的派生视图,通过流自动、按序地更新。 这样数据流向清晰、单向,避免了”A 写了 B 没同步”的混乱。批处理和流处理在这里融合——用同一份不可变的事件日志,既能流式增量更新,也能批量重跑重建。
二、拆解数据库 (Unbundling the Database)
这是全书最有想象力的观点。传统数据库把很多功能捆绑在一起:存储、索引、复制、缓存、物化视图……作者提议换个视角看整个数据生态:
把公司里所有的数据系统,看成一个巨大的、被拆开的数据库——事件日志相当于它的”预写日志/复制日志”,各个专用系统(搜索、缓存、数仓)相当于它的”不同索引/物化视图”,而连接它们的数据流,相当于数据库内部触发视图更新的机制。
也就是说,我们用组合各种专用工具(unbundling)的方式,重新搭出一个”数据库”的功能,但每个部件都能选最适合的实现。这比指望一个大而全的数据库要灵活得多。
💡 架构视野 tip:这个”拆解数据库”的比喻是本书留给你最值钱的思维模型之一。当你在真实公司里看到一堆系统(MySQL + Elasticsearch + Redis + Kafka + Snowflake……)觉得杂乱时,试着用这个视角重新组织它:谁是真相之源?数据怎么单向地从它流向各个派生视图? 一旦理清这条”数据流主干”,整个架构就从”一团乱麻”变成”一个逻辑清晰的大数据库”。这也是数据工程师最核心的架构能力。
由此引出一种新的应用设计范式——围绕数据流设计应用:像电子表格一样,当源数据变化,所有依赖它的派生数据自动更新;应用订阅变化流,而不是反复轮询查询。
三、追求正确性
在不可靠的系统上怎么保证结果正确?关键思想:
- 端到端论证 (end-to-end argument):像”恰好一次”“去重”这类正确性保证,不能只靠某一层(比如只靠 TCP 或只靠数据库)实现,往往需要贯穿从始至终的机制——例如给每个操作一个唯一 ID,用幂等来抵御重复。
- 不靠协调也能保证约束:像唯一性这类约束,可以借助基于日志的顺序来达成,而不必依赖昂贵的分布式事务/强协调。
- 信任但要核查 (trust, but verify):不要盲目相信存储永远不出错。应设计能自我审计、自我校验的系统,定期检查派生数据和真相之源是否一致(数据可能因 bug、硬件问题悄悄损坏)。这呼应了第 8 章”分布式系统里什么都不能全信”的基调。
四、做正确的事(伦理反思)
全书以一段严肃的伦理讨论收尾——技术不是价值中立的。作者提醒我们警惕:
- 预测性分析的偏见:算法可能固化甚至放大歧视(用历史数据训练 → 复制历史的不公);错误的数据或有偏的模型会给真人带来实实在在的伤害。
- 反馈回路:系统的预测会影响现实,现实又反过来强化预测,形成有害循环。
- 隐私与监控:大规模收集数据把用户置于被监视的境地;”数据是资产”的另一面是”数据是负债“——一旦泄露、滥用,代价惨重。
- 知情同意与用户自主权:用户往往无法真正理解自己的数据被如何使用。
💡 给未来工程师的 tip:作者把伦理放在全书最后,是有深意的——技术能力越强,责任越大。 你学会了设计强大的数据系统,就更要问:”这份数据该不该收集?这个模型会不会伤害到某些人?如果泄露会怎样?”这不是软性的口号,而是专业素养的一部分。一个好的数据工程师,既要让系统可靠、可扩展、可维护(回到第 1 章的三把尺子),也要让它正当。
第三部分小结
| 章 | 主题 | 一句话核心 |
|---|---|---|
| Ch10 批处理 | 处理有界的大批数据 | 源自 Unix 哲学(单一职责/统一接口/输入不可变);MapReduce→数据流引擎(Spark),靠不可变输入 + 重算容错 |
| Ch11 流处理 | 处理无界的持续事件 | 基于日志的消息(Kafka)可回放、多消费者;把数据库变化当成流(CDC/事件溯源),日志是真相、表是派生视图;小心事件时间 vs 处理时间 |
| Ch12 未来 | 把生态统一成一个大系统 | 拆解数据库:一个真相日志 + 一堆派生视图,数据单向流动;端到端保正确;最后回归伦理责任 |
贯穿第三部分(也是全书)的终极思想:现实中我们不靠一个万能数据库,而是把一个不可变的事件日志当作真相之源,让各种专用系统作为它的派生视图,通过数据流保持同步。这个”日志 + 派生视图”的世界观,把批处理、流处理、数据库、缓存、索引统一了起来。而所有这一切的最终目的,回到全书开篇的三把尺子——建造可靠、可扩展、可维护,并且正当的数据系统。
留给自己的思考题
- 我接触的系统里,有没有明确的”真相之源”?还是同一份数据在多个系统里各自为政、经常对不上?
- 我们有没有用 CDC / 事件流让派生系统(搜索、缓存、数仓)自动同步?还是靠人工/定时脚本?
- 我做过的任何”按时间统计”,用的是事件时间还是处理时间?迟到数据会不会让它出错?
- 用”拆解数据库”的视角重画我司的数据架构,数据流主干会是什么样?
- 我经手的数据,有没有想过它的伦理与隐私风险——该不该收集、泄露了会怎样?