专业书籍精读 · DDIA · 第 11 章
Designing Data-Intensive Applications · Ch 11 · Martin Kleppmann · 2017
你刷信用卡的一瞬间弹出「消费提醒」、打车软件实时挪动的那个小车、直播间蹭蹭涨的在线人数——这些都不能「等到今晚一起算」。上一章的批处理像每天洗一大桶衣服,攒够一批才跑,你看到的永远是「昨天的世界」。这一章讲的是另一种活儿:数据来一条、就处理一条,让结果永远贴着「此刻」——这就是流处理。
把每一件发生的事——一次点击、一笔付款、一台传感器的读数——都想成一张写好就不改的小便利贴:「几点几分,谁,干了什么」。流处理就是在便利贴不断飞来的传送带旁边站个人,来一张看一张、随手就把账记上、把该报的警报出了,而不是等一麻袋攒满再拆。
批处理有个改不掉的毛病:它总在等——等数据攒够一批、等整批算完,你拿到的永远是几小时前的旧账。可现实里一大堆事等不起:盗刷要当场拦、大促库存要实时看、故障要立刻报。你总不能为了看一眼此刻,把一整天的数据从头再算一遍。
撑起流处理的有两个朴素点子。其一,一条永不删的流水账。别把便利贴看完就扔——把它们一张接一张钉在一条长长的传送带上、永不撕掉;哪个部门要看,就自己带个书签、读到哪记到哪,随时能把书签往回拨、重读一遍(这正是 Kafka 干的事)。其二,账本与余额是一回事。你的银行余额,无非是把所有流水从头加一遍的结果;反过来,只要把每一笔流水都留着,任何时刻的余额都能重新算出来。「存下每一次变化、而不只存最新状态」——这个小小的调头,是整章的魂。
有了这条可回放的流水账,一处发生的变化能秒级扇给所有需要它的人:数据库一改,搜索索引、缓存、报表跟着实时更新;风控在你手指离开屏幕前就算完了要不要拦。要「立等可取」的场景就上流,能「今晚跑明早出」的就交给批。一句诚实的代价:一旦追求实时,「时间」本身就变棘手了——地铁里发出的消息可能晚几分钟才到,事件乱序、迟到,让「这一分钟到底发生了什么」永远没有一条铁定的截止线,只能划一道大概的界、承认偶有漏网。
流处理 = 数据来一条算一条,让结果永远贴着此刻,图的是「立等可取」。两根支柱:一条永不删、可回放的流水账(Kafka),和「存每一次变化、而非只存最新状态」——余额不过是流水加总。代价是:实时让「时间」变滑,乱序与迟到的事件没有铁定截止线。
想进到 CDC、流表二象性、窗口与水位、exactly-once 和示意图? → 切到精读版
这一章讲流处理(stream processing):把数据看成一条永不结束的事件流(unbounded stream),事件一到就处理,让派生结果(索引、缓存、实时指标、告警)持续贴近「此刻」。它和上一章批处理共用「不可变数据 + 派生」的世界观,只是把「有界的一大批」换成「无界的连续流」。两根支柱撑起全章:日志型消息代理(log-based message broker,如 Kafka)——一条持久、有序、可回放的事件日志;事件溯源 / 流表二象性——把「每一次变化」而非「最新状态」当作事实之源。一句反直觉:你以为流处理的难点是「算得快」,其实是「时间」——事件时间与处理时间的错位,才是这一章真正的深水区。
(用户A, 12:03, 点了商品X)。流处理全程在搬运事件。本章是全书 Part III「派生数据(Derived Data)」的第二块拼图,紧接上一章批处理(Ch10)。Ch10 处理有界、静态的一大批数据;本章掉头处理无界、实时的连续流。二者不是对手而是同一世界观的两副面孔——都把数据当不可变输入、从中派生结果,区别只在「批 vs 流」。它上承 Ch10 的批处理、也回收了 Ch5 复制里「复制日志」、Ch7「事务与 exactly-once」的伏笔;下启终章 Ch12——作者要用流的思路把整个数据系统重新组织。对应现实里的 Kafka / Flink / Spark Streaming 生态、实时数仓、风控与监控管线。
假设你在做一个电商风控系统:每秒涌进 50,000 笔交易,你要在用户点下「支付」到银行返回之间——不到 200 毫秒——判断这笔是不是盗刷。再假设运营要一块实时大屏,看当下每分钟的 GMV、各城市的下单热力。这类活的共同点是:结果的价值随时间飞快衰减,「今晚跑、明早出」的批处理彻底没用——等你算完,钱早被盗刷走了、大促也结束了。
可你又不能把批处理那套「每小时把全量从头算一遍」硬缩短到「每秒算一遍」——那会把机器和数据库压垮。真正要解决的是三个环环相扣的难题:① 怎么把源源不断的事件可靠、有序地从产生方送到处理方,还能在处理方崩溃重启后从断点续上、甚至回放历史;② 怎么让一处数据的变化实时、正确地传导到所有派生系统(搜索、缓存、报表),不让它们各说各话;③ 怎么在事件乱序、迟到的现实里,对「一个时间窗口」给出可信的聚合。不解决这些,你要么忍几小时延迟,要么写一堆脆弱的、丢消息又对不上账的实时脚本。
流处理的起点是怎么把事件从生产者送到消费者。最朴素的做法是让消费者不停轮询数据库看有没有新数据,但太笨。专门的消息代理(message broker)应运而生,历史上分两大流派——理解它们的差别,是理解 Kafka 为什么重要的关键:
日志型 broker 的意义在于:它同时具备数据库的持久性和消息队列的低延迟——把「传递消息」和「存储数据」合成一件事。这正是让「流」能当作可信事实之源的技术地基。代价是:一个分区只能被一个消费者顺序读,想让 100 个消费者并行分摊同一批消息、且乱序无所谓时,传统队列反而更灵活。
把「流」这个抽象接到数据库上,就打通了本章最有用的一类应用:让所有系统保持同步。现实里你几乎不可能只用一个数据库——主库之外还有搜索引擎(Elasticsearch)、缓存(Redis)、数仓、推荐特征库,它们都要跟着主库的数据变。传统做法是应用层双写(改完数据库再手动改搜索),但双写极易在并发或部分失败下对不上账。更干净的办法是把数据库的变化变成一条流:
变更数据捕获(CDC):主库每发生一次增删改,就抽成一条事件(通常直接读它的复制日志 / WAL),推进一条流。所有下游订阅这条流、按序应用,就能最终和主库一模一样——本质上是把 Ch5 的「主从复制」推广到异构系统之间。工具如 Debezium 就干这个。
事件溯源(event sourcing):更彻底的一步——干脆不把「当前状态」当事实之源,而把「导致状态的每一个事件」当事实之源。购物车不存「现在有 2 件商品」,而存下 加入X、加入Y、移除X 这串不可变事件,当前状态由重放算得。好处:完整审计(每步留痕、可追责)、可回到任意历史时刻、同一份事件流能派生多种视图(同样一串下单事件,既算库存、又算营收、还喂推荐)。
CDC 和事件溯源背后是同一个深刻观察,作者称之为流与表的二象性(stream-table duality):
这不是文字游戏,而是一种架构世界观:把数据库「里外翻转」——传统数据库把「变更日志」藏在内部、只把「当前状态」露给你;流处理反过来,把变更日志本身升格为一等公民(就放在 Kafka 里),让『当前状态』退化成随时可从日志重建的、多份的物化视图。于是搜索索引、缓存、报表全成了「同一条事实日志的不同投影」,天然一致、坏了能从日志重建。这正是 Kleppmann 那句著名的「turning the database inside out(把数据库里外翻转)」,也是终章 Ch12 的思想内核。
这是本章公认最烧脑、也最见真章的一节。批处理里「时间」不是问题(数据都在,随便算);到了流处理,你想算「过去一分钟有多少笔交易」,立刻撞上一个哲学难题:哪个「一分钟」?
二者常对不上:用户在地铁里下了单,手机没信号,5 分钟后出站才把事件补发上来。若按处理时间算,这笔会被错记进 5 分钟后的那个窗口。按处理时间算简单但结果会因系统忙闲、网络抖动而失真;按事件时间算才「对」,但你会遇到一个死结:你永远不能确定某个时间窗口的事件是不是都到齐了——总可能还有个掉队的(straggler)在路上。
解药是水位线(watermark):系统维护一条启发式的界线「事件时间 T 之前的,我认为基本都到了」,一旦水位越过窗口的结束边界,就先把这个窗口的结果算出来发下去。之后再来的迟到事件怎么办?三种态度:直接丢弃(并计数报警)、发一条更正(撤回旧结果重发)、或允许一段宽限期再迟到就丢。这里没有免费午餐——想等得久一点、少漏一个,就得牺牲延迟;想快点出结果,就得容忍偶有漏网。
切窗口本身也有讲究,常见四种:滚动窗口(tumbling)——首尾相接不重叠(「每整分钟」);跳动窗口(hopping)——固定长度但按更小步长滑动、会重叠(「每分钟统计过去 5 分钟」);滑动窗口(sliding)——任意时刻都看「往前 N 分钟」;会话窗口(session)——没有固定长度,按「同一用户一连串活动、中间静默超过若干分钟才算断」来切。
和批处理一样,流上也常要 join,但因为数据在「流动」,形态更微妙,分三类:①流-流 join(窗口连接)——把两条流在一个时间窗口内配对,如「搜索事件」配「点击事件」算点击率,难点是两边都在动、要缓存一个窗口的状态等对方;②流-表 join(补全 / enrichment)——用一张(缓慢变化的)表去丰富流里的每条事件,如给每条点击事件补上用户画像,实现上常把表通过 CDC 也变成流、在本地维护一份副本;③表-表 join(维护物化视图)——两条 changelog 流合并出一个持续更新的联表视图。
容错与 exactly-once(恰好一次)是流处理的终极考题。批处理失败大不了整批重跑(输入还在);流是无界的、没有「跑完」,崩溃重启后怎么保证「每个事件的效果恰好生效一次、不重不漏」?三条主流路线:①微批(microbatching,Spark Streaming)——把流切成一秒一个的小批,退化成批处理来容错,简单但延迟受限于批长;②检查点 + 幂等(checkpoint,Flink)——周期性给算子状态打快照,崩溃后回滚到上个检查点重放,配合幂等写出,达成「效果恰好一次」;③原子提交——把「推进读进度」和「写结果」绑进一个事务一起成败(如 Kafka 事务)。注意:exactly-once 指「对外部效果恰好生效一次」,而非「物理上只处理一次」——底层往往仍靠重放 + 去重 / 幂等实现。
流处理的选型,核心是在可回放性、顺序、实时性、正确性代价之间权衡。三张表把最常考、最实用的对比拎出来。
表 1 · 传统消息队列 vs 日志型消息代理(Kafka)
| 传统队列(AMQP/JMS) | 日志型 broker(Kafka) | |
|---|---|---|
| 消费后 | 签收即删除,一次性 | 只追加、读过不删,按时间/容量保留 |
| 能否回放 | 不能——消费完就没了 | 能——拨回 offset 即重读历史 |
| 顺序 | 并发消费下容易乱序 | 分区内严格保序 |
| 扩展方式 | 多消费者分摊一个队列(负载均衡) | 按 key 分区,分区内单消费者、跨分区并行 |
| 最擅长 | 任务分发(每条只需被处理一次,如发邮件) | 事件流 / 数据同步 / 需回放与多下游 |
表 2 · 事件时间 vs 处理时间(及迟到事件对策)
| 按处理时间 | 按事件时间 | |
|---|---|---|
| 含义 | 事件被处理到的时刻 | 事件真正发生的时刻 |
| 实现难度 | 简单,无需等待 | 难,需水位线 + 缓存状态 |
| 结果正确性 | 会因系统忙闲/网络抖动失真 | 语义正确,反映真实业务时刻 |
| 迟到事件 | 被错记进后面的窗口 | 过水位后到达 → 丢弃 / 更正 / 宽限期 |
| 何时选 | 只关心「系统吞吐节奏」的粗略监控 | 要「真实业务时刻」的账(计费、风控、分析) |
表 3 · 流上三种 join 怎么选
| 类型 | 做什么 | 典型场景 | 难点 / 代价 |
|---|---|---|---|
| 流-流 (窗口) | 两条流在时间窗口内配对 | 搜索事件 × 点击事件算点击率 | 两边都在动,需缓存整个窗口的状态 |
| 流-表 (补全) | 用一张表丰富流中每条事件 | 给点击补上用户画像 | 需本地维护表副本(常靠 CDC 同步) |
| 表-表 (物化视图) | 两条 changelog 合出持续更新的联表 | 时刻维护「用户+订单」联合视图 | 状态大;表变化也要及时传导 |
流处理是当今实时数据栈的骨架。Apache Kafka(LinkedIn 2011 年开源)把「日志」抬成了基础设施,几乎成了事件流的事实标准;其上 Kafka Streams、Apache Flink、Spark Streaming、Storm、Samza 层层生长。CDC 工具(Debezium)把老数据库接进流;实时数仓、实时风控、推荐特征、监控告警、物联网遥测,骨子里都是这一章的思想。作者提出的「把数据库里外翻转」(unbundling / 流表二象性)更直接定义了 Ch12 的方向。面试里,「Kafka 为什么能既持久又低延迟」「exactly-once 到底怎么实现」「事件时间和处理时间的区别与水位线」几乎是数据 / 后端岗高频题——答案全在本章。
① 一句话:流处理 = 把数据当无界事件流、来一条算一条,让派生结果持续贴近此刻,图「立等可取」。
② 与批处理是一枚硬币两面:同样「不可变输入 + 派生」,区别只在有界的一批 vs 无界的流。
③ 支柱一:日志型消息代理(Kafka)——只追加、有序、读过不删,靠 offset 书签可回放;比「签收即删」的传统队列多了持久与重放。
④ 支柱二:CDC + 事件溯源——把「每一次变化」而非「最新状态」当事实之源,用一条流让搜索/缓存/数仓与主库最终一致,取代脆弱的双写。
⑤ 灵魂:流表二象性——表是流的快照(重放加总),流是表的变更日志;「把数据库里外翻转」,让状态成为可从日志重建的物化视图。
⑥ 最烧脑的是时间:事件时间 vs 处理时间;用窗口 + 水位线决定「一个时间窗何时收工」,迟到事件只能丢 / 更正 / 宽限,无两全。
⑦ 三种 join(流-流窗口、流-表补全、表-表物化视图);容错的 exactly-once 指「效果恰好一次」,靠微批 / 检查点+幂等 / 原子提交实现。
⑧ 落地:Kafka → Flink / Spark Streaming / Kafka Streams,撑起实时数仓、风控、监控;Kreps「The Log」与 Google Dataflow Model 是思想源头,直通终章 Ch12。