13 流式系统的哲学

如果一件事物以另一件事物为目的,那么它的终极目的就不可能只是保存自身。因此,船长不会把保全受托的船只当作终极目的,因为船只另有其目的,也就是航行。
(这段话常被引述为:如果船长的最高目标是保护船只,他就会让船永远停在港口。)
圣托马斯・阿奎那,《神学大全》(1265—1274)
我们在 第 2 章 中讨论了构建 可靠(reliable)、可伸缩(scalable)、可维护(maintainable)的应用与系统这一目标。这些主题贯穿全书各章:例如,我们讨论了许多有助于提升可靠性的容错算法、提升可伸缩性的分片,以及提升可维护性的演化与抽象机制。
在本章中,我们将把所有这些想法汇集起来,并特别以 第 12 章 的 流处理(stream processing)与 事件驱动架构(event-driven architecture)为基础,形成一套能够实现上述目标的应用开发哲学。与前几章相比,本章的立场更为鲜明:它将深入阐述一种特定的哲学,而不是比较多种不同的方法。
数据集成
本书反复出现的一个主题是:对于任何给定的问题,往往都有若干种解决方案,而每种方案都有不同的优点、缺点与利弊权衡。例如,在 第 4 章 讨论存储引擎时,我们看到了日志结构存储、B 树和列式存储;在 第 6 章 讨论复制时,我们看到了单主、多主和无主复制。
如果你面临的问题是“我想存储一些数据,稍后再把它查出来”,那么并不存在唯一正确的解决方案;不同的方法各自适用于不同的情形。软件实现通常必须选择一种具体的方法。要让一条代码路径既稳健又有良好的性能,本就已经很难;试图在一个软件中包办一切,几乎注定会得到糟糕的实现。
因此,选择哪一种软件工具最合适,同样取决于具体情形。每一种软件——即便是所谓的“通用”数据库——都是针对某种特定的使用模式设计的。
面对如此繁多的选择,第一个挑战便是弄清各种软件产品分别适合什么情形。供应商不愿告诉你自己的软件不适合哪些工作负载,这完全可以理解;不过,希望前面的章节已经让你知道应该提出哪些问题,从而读懂言外之意,更好地理解其中的权衡。
然而,即使你已经完全掌握了工具及其适用情形之间的对应关系,仍然会遇到另一个挑战:在复杂应用中,数据往往会以许多不同的方式使用。几乎不可能有一种软件适合数据的 所有 使用情形,因此,为了提供应用所需的功能,你不可避免地需要把几种不同的软件组合起来。
通过衍生数据组合专用工具
例如,为了支持任意关键词查询,通常需要把 OLTP 数据库与全文检索索引集成起来。尽管某些数据库(例如 PostgreSQL)内置的全文索引功能足以满足简单应用的需要 1,但更复杂的检索功能仍然需要专业的信息检索工具。反过来,搜索索引通常又不适合作为持久的权威记录系统,因此许多应用都需要组合使用这两种工具,才能满足全部需求。
我们在 “保持系统同步” 中谈到过数据系统的集成问题。数据表示越多,集成就越困难。除了数据库和搜索索引,你可能还需要在分析系统中保存数据副本,例如数据仓库、批处理(batch processing)系统或流处理系统;维护由原始数据衍生而来的缓存或反规范化对象;让数据经过机器学习、分类、排名或推荐系统;或者根据数据变更发送通知。
理解数据流
为了满足不同的访问模式,如果同一份数据的副本需要保存在多个存储系统中,你就必须清楚地界定输入和输出:数据最先写到哪里?哪些表示是从哪些来源衍生出来的?怎样才能以正确的格式,把数据送到所有正确的位置?
例如,你可以让数据首先写入作为权威记录系统的数据库,捕获该数据库中的变更(参见 “变更数据捕获”),再按相同顺序把这些变更应用到搜索索引。如果变更数据捕获(CDC)是更新索引的唯一途径,你就可以确信索引完全衍生自权威记录系统,因而与之保持一致(软件缺陷除外)。在这个系统中,只有写入数据库才能提供新的输入。
如果允许应用同时直接写入搜索索引和数据库,就会引入 图 12-4 所示的问题:两个客户端并发提交相互冲突的写入,而两个存储系统却以不同的顺序处理它们。此时,无论数据库还是搜索索引,都不能“说了算”来决定写入顺序;它们可能作出相互矛盾的决定,并从此永久不一致。
如果能够把所有用户输入都汇入一个系统,由它决定全部写入的顺序,那么只需按相同顺序处理这些写入,就能更容易地衍生出数据的其他表示。这正是我们在 “共识的实践” 中见过的状态机复制方法的一种应用。使用变更数据捕获还是事件溯源日志并不是最重要的,真正重要的是先确定一个全序。
根据事件日志更新衍生数据系统,往往可以做到 确定性(determinism)与 幂等性(idempotence;参见 “幂等性”),因而很容易从故障中恢复。
衍生数据与分布式事务
让不同数据系统彼此保持一致的经典方法是使用分布式事务,如 “两阶段提交(2PC)” 所述。相比之下,使用衍生数据系统的效果如何?
从抽象层面看,两者以不同手段实现了相似的目标。分布式事务使用锁来实现互斥,以此决定写入顺序;而 CDC 与事件溯源则用日志来排序。分布式事务通过原子提交确保变更恰好生效一次;基于日志的系统通常依靠确定性重试与幂等性。
两者最大的区别在于:事务系统通常保证一个值写入后,立即就能读到它的最新值(参见 “读己之写”)。衍生数据系统则往往异步更新,因此默认并不保证读操作所见的数据是最新的。
在一些范围有限、愿意承担分布式事务成本的环境中,分布式事务已得到成功应用。然而,XA 的容错性和性能都不理想(参见 “跨不同系统的分布式事务”),这严重限制了它的实用性。或许可以设计出一种更好的分布式事务协议,但要让这种协议获得广泛采用,并与现有工具集成,将极具挑战,短期内也不太可能实现。
由于目前还没有一种得到广泛支持的优良分布式事务协议,基于日志的衍生数据是集成不同数据系统最有前景的方法。不过,读己之写等保证确实很有用;一味告诉所有人“最终一致性不可避免——接受现实,学会应付吧”,并无助益(至少在没有妥善说明 如何 应对时是这样)。
本章稍后将讨论一些在异步衍生系统之上实现更强保证的方法,努力在分布式事务与基于日志的异步系统之间找到一条中间道路。
全序的局限
对于规模足够小的系统,构建全序事件日志完全可行(采用单主复制的数据库如此流行,便证明了这一点:它们构建的正是这样的日志)。不过,随着系统扩展到更大、更复杂的工作负载,局限便开始显现:
在大多数情况下,构建全序日志要求所有事件都经过 单个领导者节点,由它决定顺序。如果事件吞吐量超出单台机器的处理能力,就需要把日志分片到多台机器上。此时,两个不同分片中的事件孰先孰后便不再明确。
如果服务器分布在多个 地理上分散 的区域——例如,为了容忍整个数据中心离线——通常会在每个数据中心分别设置领导者,因为网络延迟会使跨数据中心的同步协调效率低下。这意味着,来自两个不同数据中心的事件之间没有确定的顺序。
当应用以 微服务 形式部署时,一种常见的设计选择是把每项服务及其持久状态作为独立单元部署,不让服务之间共享持久状态。当两个事件分别产生于不同服务时,它们之间没有确定的顺序。
一些应用会在客户端维护状态:用户输入后立即更新,不等待服务器确认,甚至在离线时仍能继续工作。在这样的应用中,客户端与服务器很可能以不同的顺序看到事件。
从形式上说,决定事件全序的问题称为 全序广播(total order broadcast),它等价于共识(参见 “共识的多面性”)。大多数共识算法针对的是单个节点吞吐量足以处理整个事件流的情形,并没有提供让多个节点共同分担事件排序工作的机制。
排序事件以捕获因果关系
如果事件之间没有因果联系,缺少全序并不是什么大问题,因为并发事件可以按任意顺序排列。有些情形也很容易处理:例如,同一个对象有多次更新时,只要把特定对象 ID 的所有更新都路由到同一个日志分片,就能为它们建立全序。然而,因果依赖有时会以更隐蔽的方式出现。
例如,考虑一个社交网络服务,以及两个曾是情侣、但刚刚分手的用户。其中一人先把另一人移出好友列表,然后向剩余好友发送消息,抱怨自己的前任。这个用户的本意是,不让前任看到这条粗鲁的消息,因为消息是在好友关系解除之后发送的。
然而,如果一个系统把好友关系状态和消息存放在不同的位置,那么 解除好友 事件与 发送消息 事件之间的顺序依赖可能会丢失。如果没有捕捉到这种因果依赖,负责发送新消息通知的服务就可能先处理 发送消息 事件、后处理 解除好友 事件,从而错误地向前任发出通知。
在这个例子中,通知实际上是消息与好友列表之间的连接,因此它涉及我们之前讨论过的连接时序问题(参见 “连接的时间依赖性”)。遗憾的是,这个问题似乎没有简单的答案 2、3。一些可能的起点包括:
逻辑时间戳无需协调即可提供全序(参见 “ID 生成器和逻辑时钟”),因而在全序广播不可行时或许能派上用场。然而,接收方仍需处理乱序到达的事件,而且系统还必须传递额外的元数据。
如果能用一条日志事件记录用户作出决定前所看到的系统状态,并为该事件赋予唯一标识符,那么后续事件便可引用这个事件 ID,从而记录因果依赖 4。
冲突解决算法(参见 “自动解决冲突”)有助于处理以意外顺序到达的事件。它们适合用于维护状态;但如果某个操作会产生外部副作用(例如向用户发送通知),这些算法就无能为力了。
或许未来会出现新的应用开发模式,能够高效捕捉因果依赖、正确维护衍生状态,而不必强迫所有事件都经过全序广播这一瓶颈。
批处理与流处理
数据集成(data integration)的目标,是确保数据以正确的形式出现在所有正确的位置。为此,需要消费输入,进行转换、连接、过滤、聚合、模型训练与评估,最终再写入适当的输出。批处理器和流处理器正是实现这一目标的工具。批处理与流处理的输出是衍生数据集,例如搜索索引、物化视图、向用户展示的推荐结果、聚合指标,等等。
正如我们在 第 11 章 和 第 12 章 中看到的,批处理与流处理有许多共同原则;两者最根本的区别在于,流处理器处理的是无界数据集,而批处理的输入大小已知且有限。
维护衍生状态
批处理带有很强的函数式风格(即使代码并不是用函数式编程语言编写的):它鼓励使用确定性的 纯函数(pure function),其输出只取决于输入,除了明确的输出之外没有其他 副作用(side effect);输入被视为不可变,输出则只能追加。流处理与之类似,不过它扩展了算子,使其能够维护受管理且容错的状态。
输入和输出定义明确的确定性函数,不仅有利于容错,也能简化对组织内部数据流的推理 5。无论衍生数据是搜索索引、统计模型还是缓存,都可以把它看成数据管道:从一项事物衍生出另一项事物,让一个系统的状态变更经过函数式应用代码,再把相应效果施加到衍生系统上。这样思考大有裨益。
原则上,衍生数据系统也可以同步维护,就像关系数据库在写入被索引表的同一事务中,同步更新二级索引一样。然而,异步恰恰是基于事件日志的系统能够保持稳健的原因:系统某一部分的故障可以被局限在本地;分布式事务则会在任何一个参与者失败时中止,因而容易把故障扩散到系统其他部分,将其放大。
我们在 “分片与二级索引” 中看到,二级索引常常跨越分片边界。带有二级索引的分片系统,要么需要把写入发送到多个分片(按词项分区的索引),要么需要把读取发送到所有分片(按文档分区的索引)。如果异步维护索引,这种跨分片通信也最为可靠、最具可伸缩性 6。
为应用演化而重新处理数据
维护衍生数据时,批处理和流处理都很有用。流处理可以用很低的延迟把输入中的变更反映到衍生视图中;批处理则可以重新处理大量累积的历史数据,从已有数据集衍生出新的视图。
尤其是,重新处理现有数据为维护和演化系统、支持新功能与变化后的需求,提供了一种良好机制。如果不能重新处理,模式演化就只能局限于简单改动,例如为记录添加一个新的可选字段,或者增加一种新的记录类型。而有了重新处理,就可以把数据集重组为完全不同的模型,从而更好地满足新需求。
大规模的“模式迁移”也会发生在非计算机系统中。例如,19 世纪英国铁路建设初期,轨距(两条铁轨之间的距离)存在多种相互竞争的标准。为一种轨距建造的列车无法在另一种轨距的轨道上行驶,这限制了铁路网络可能实现的互联 7。
1846 年终于确定了统一的标准轨距,其他轨距的线路都必须改造——可怎样才能在不让铁路停运几个月乃至几年的情况下完成改造?解决办法是先增加第三条铁轨,把线路改成 双轨距 或 混合轨距。这种改造可以逐步进行;完成之后,两种轨距的列车都能在同一条线路上行驶,各自使用三条铁轨中的两条。最终,所有列车都改用标准轨距后,非标准轨距所用的那条铁轨便可拆除。
以这种方式“重新处理”既有轨道,并让新旧版本并行存在,就能在数年时间里逐步改变轨距。尽管如此,这仍是一项昂贵的工程,所以非标准轨距至今依然存在。例如,旧金山湾区的 BART 系统使用的轨距就不同于美国大多数铁路。
衍生视图允许系统 渐进式 演化。如果要重组一个数据集,并不需要突然切换完成迁移。你可以把旧模式和新模式作为底层同一份数据上的两个独立衍生视图,并排维护。随后,先把少量用户转向新视图,以测试其性能并发现缺陷,同时让大多数用户继续使用旧视图。之后逐步提高访问新视图的用户比例,最终删除旧视图 8、9。
这种渐进迁移的妙处在于:如果出了问题,过程中的每个阶段都很容易逆转,始终有一个可用的系统供你退回。由于不可逆损害的风险降低,你可以更有信心地继续推进,从而更快地改进系统 10。
统一批处理与流处理
统一批处理与流处理的一项早期提案是 Lambda 架构 11。它存在不少问题 12,如今已很少使用。较新的系统允许在同一系统中实现批计算(重新处理历史数据)与流计算(在事件到达时处理事件)13;这种方法有时称为 Kappa 架构 12。
要在一个系统中统一批处理与流处理,需要具备以下功能:
能够让历史事件通过处理近期事件流的同一个处理引擎重放。例如,基于日志的消息代理可以重放消息,一些流处理器也能从分布式文件系统或对象存储中读取输入。
为流处理器提供 恰好一次语义(exactly-once semantics)——也就是说,即使实际发生了故障,也要确保输出与从未发生故障时相同。与批处理一样,这要求丢弃所有失败任务的部分输出。
提供按事件时间而非处理时间划分窗口的工具,因为重新处理历史事件时,处理时间没有任何意义。例如,Apache Beam 提供了表达这类计算的 API,之后可以用 Apache Flink 或 Google Cloud Dataflow 来运行。
分拆数据库
在最抽象的层面上,数据库、批处理器、流处理器和操作系统执行着相同的功能:存储一些数据,并允许你处理和查询这些数据 14、15。数据库把数据存储为某种数据模型中的记录(表中的行、文档、图中的顶点等),操作系统的文件系统则把数据存储在文件中——但二者本质上都是“信息管理”系统 16。正如我们在 第 11 章 中看到的,批处理器就像 Unix 的分布式版本。
当然,实际差异仍然很多。例如,许多文件系统无法很好地处理一个包含 1000 万个小文件的目录,而数据库中有 1000 万条小记录却完全稀松平常。尽管如此,操作系统与数据库之间的异同仍然值得探究。
Unix 与关系数据库采用截然不同的哲学来处理信息管理问题。Unix 认为自己的目标是向程序员提供一种合乎逻辑、但层次相当低的硬件抽象;关系数据库则希望向应用程序员提供高层次抽象,隐藏磁盘数据结构、并发、崩溃恢复等复杂问题。Unix 发展出了管道以及本质上只是字节序列的文件,数据库则发展出了 SQL 与事务。
哪种方法更好?当然要看你想要什么。Unix 的“简单”,在于它只是硬件资源外面一层相当薄的包装;关系数据库的“简单”,则在于一条简短的声明式查询就能借助大量强大的基础设施(查询优化、索引、连接方法、并发控制、复制等),而查询作者无需了解实现细节。
这两种哲学之间的张力已经持续数十年(Unix 和关系模型都出现于 20 世纪 70 年代初),至今仍未消解。例如,可以把 NoSQL 运动理解为:试图把 Unix 式的低层抽象方法应用于分布式 OLTP 数据存储领域。
本节试图调和这两种哲学,希望能够兼取二者之长。
组合使用数据存储技术
在本书的过程中,我们讨论了数据库提供的各种功能及其工作原理,其中包括:
二级索引,使你可以根据字段值高效搜索记录;
物化视图,一种预先计算的查询结果缓存;
复制日志,使其他节点上的数据副本保持最新;以及
全文检索索引,允许在文本中进行关键词搜索,一些关系数据库也内置了这种索引 1。
在 第 11 章 和 第 12 章 中,我们也遇到了类似主题。我们谈过如何构建全文检索索引、如何维护物化视图,以及如何通过变更数据捕获,把数据库中的变更复制到衍生数据系统。
数据库内置的功能,与人们使用批处理器和流处理器构建的衍生数据系统,似乎存在诸多相似之处。
创建索引
想一想,在关系数据库中运行 CREATE INDEX 创建新索引时会发生什么。数据库必须扫描表的一致性快照,取出所有要索引的字段值,对其排序,再写出索引。接着,它还必须处理拍摄一致性快照后积压的所有写入(假设创建索引时没有锁住表,因此写入可以继续)。完成之后,每当事务写入该表,数据库都必须继续更新索引。
这个过程与设置新的从库副本极其相似(参见 “设置新的副本”),也很像在流处理系统中引导变更数据捕获(参见 “初始快照”)。
每次运行 CREATE INDEX,数据库本质上都是在重新处理现有数据集,并把索引衍生为既有数据之上的一个新视图。现有数据可能是状态快照,而不是曾经发生过的所有变更的日志,但二者密切相关。
一切的元数据库
从这个角度看,整个组织的数据流开始像一个巨大的数据库 5。每当批处理、流处理或 ETL 过程把数据从一种位置和形式传送到另一种位置和形式,它就像是在充当维护索引或物化视图的数据库子系统。
这样看来,批处理器和流处理器就像触发器、存储过程与物化视图维护算法的精巧实现;它们维护的衍生数据系统则像不同种类的索引。例如,关系数据库可能支持 B 树索引、哈希索引、空间索引以及其他索引。在新兴的衍生数据系统架构中,这些能力不再作为单个集成数据库产品的功能来实现,而是由各种不同的软件提供,运行在不同机器上,并由不同团队管理。
这些发展未来会把我们带向何方?如果从“没有任何一种数据模型或存储格式适合所有访问模式”这一前提出发,那么仍有两条途径,可以把不同的存储与处理工具组合成一个协调一致的系统:
- 联合数据库:统一读取
可以为各种底层存储引擎和处理方法提供统一的查询接口——这种方法称为 联合数据库(federated database)或 多模存储(polystore)17、18。例如,PostgreSQL 的 外部数据封装器(foreign data wrapper)功能符合这种模式,Trino、Hoptimator 和 Xorq 等联合查询引擎也是如此。需要专用数据模型或查询接口的应用仍可直接访问底层存储引擎;而希望组合不同位置数据的用户,则可以通过联合接口轻松完成。
联合查询接口延续了关系模型的传统:提供带有高层查询语言和优雅语义的单一集成系统,但其实现非常复杂。
- 分拆数据库:统一写入
联合能够解决跨多个不同系统的只读查询问题,却无法很好地解决这些系统之间的写入 同步。前面说过,在单个数据库中,创建一致的索引是一项内置功能。当我们组合多个存储系统时,同样需要确保所有数据变更最终到达所有正确的位置,即使发生故障也不例外。让存储系统更容易可靠地连接起来(例如通过变更数据捕获和事件日志),就像把数据库的索引维护功能 分拆 出来,使它能够跨不同技术同步写入 5、19。
分拆方法遵循 Unix 的传统:小工具各自做好一件事 20,通过统一的低层 API(管道)通信,再用更高层的语言(shell)组合起来 14。
让分拆行得通
联合与分拆是一枚硬币的两面:都是用不同组件组成可靠、可伸缩、可维护的系统。联合只读查询需要把一种数据模型映射到另一种数据模型,这需要一些思考,但归根结底是个相当容易处理的问题。让多个存储系统的写入保持同步,才是更困难的工程问题,因此这里将重点讨论它。
同步写入的传统方法要求跨异构存储系统使用分布式事务 17,但如前所述,这种方法问题重重。单个存储系统或流处理系统内部的事务是可行的;然而,当数据跨越不同技术之间的边界时,带有幂等写入的异步事件日志要稳健、实用得多。
例如,一些流处理器内部使用分布式事务来实现恰好一次语义,而且效果可以很好。然而,如果一项事务需要涉及由不同团队编写的系统(例如从流处理器把数据写入分布式键值存储或搜索索引),缺少标准化事务协议就会使集成难上加难。带有幂等消费者的有序事件日志是一种简单得多的抽象,因此更有可能跨异构系统实现 5。
基于日志的集成有一项巨大优势:各个组件之间 松散耦合(loose coupling)。这种优势体现在两个方面:
在系统层面,异步事件流使整个系统更能抵御个别组件中断或性能下降。如果某个消费者速度很慢或发生故障,事件日志可以缓冲消息,让生产者和其他消费者不受影响地继续运行。故障消费者修复后可以追赶进度,因此不会漏掉任何数据,而故障也被限制在局部。相比之下,分布式事务中的同步交互往往会把局部故障升级为大规模失效。
在人员层面,分拆数据系统使不同团队可以彼此独立地开发、改进和维护不同的软件组件与服务。专业化让每个团队都能专心做好一件事,并通过定义明确的接口与其他团队的系统交互。事件日志提供的接口既足够强大,可以表达相当强的一致性属性(得益于事件的持久性与顺序),又足够通用,几乎适用于任何种类的数据。
分拆式系统与集成式系统
即使分拆确实成为未来的方向,它也不会取代现有形态的数据库——人们仍会一如既往地需要数据库。流处理器需要数据库来维护状态,批处理器与流处理器的输出也需要数据库来提供查询服务。专用查询引擎对特定工作负载仍然十分重要:例如,数据仓库中的查询引擎针对探索式分析查询进行了优化,很擅长处理这类工作负载。
运行多种基础设施所带来的复杂性可能是个问题:每种软件都有学习曲线、配置问题和运维怪癖,因此值得尽量减少部署中的活动部件。对于其设计所针对的工作负载,与用应用代码把多个工具拼接起来的系统相比,单一集成软件产品也可能提供更好、更可预测的性能 21。为并不需要的规模构建系统只是徒费力气,还可能把你锁定在僵化的设计中;这实际上是一种过早优化。
分拆的目标并不是在特定工作负载的性能上与单个数据库竞争;它的目标是让你能够组合多个不同的数据库,从而在远比单一软件所能覆盖的工作负载范围内取得良好性能。它追求的是广度,而不是深度。
因此,如果有一种技术能够满足你的全部需求,最好的选择很可能就是直接使用该产品,而不是试图用更低层的组件自行重新实现。只有当没有任何单一软件能够满足全部需求时,分拆与组合的优势才会显现。
用于组合数据系统的工具正在不断完善:Debezium 可以从许多数据库中抽取变更流;Kafka 协议正在成为事件流事实上的标准;增量视图维护引擎(参见 “增量视图维护”)则让复杂查询的缓存可以预先计算并持续更新。
围绕数据流设计应用
底层数据一旦变化便更新衍生数据,这个基本想法并不新鲜。例如,电子表格早已有强大的数据流编程能力 22:你可以在一个单元格中放入公式(比如求另一列单元格之和),每当公式的任何输入发生变化,公式结果都会自动重新计算。这正是我们希望数据系统在系统层面做到的事:数据库中的一条记录发生变化时,该记录的所有索引都应自动更新,所有依赖它的缓存视图或聚合结果也应自动刷新。你不必操心刷新过程的技术细节,只需要相信它会正确运行。
因此,大多数数据系统仍有许多地方需要向 1979 年的 VisiCalc 学习 23。与电子表格不同,当今的数据系统必须容错、可伸缩,并能持久地存储数据;它们还必须能够集成不同团队在不同时间编写的异构技术,并复用既有库与服务。指望所有软件都使用某一种语言、框架或工具开发,并不现实。
本节将进一步展开这些想法,探讨如何围绕数据库分拆与数据流的理念来构建应用。
应用代码作为衍生函数
一个数据集从另一个数据集衍生出来时,需要经过某种转换函数。例如:
二级索引是一种衍生数据集,其转换函数很直接:对于基础表中的每一行或每一份文档,取出要索引的列值或字段值,再按这些值排序(假设使用按键排序的 SSTable 或 B 树索引)。
全文检索索引的创建过程,是先应用语言检测、分词、词干提取或词形还原、拼写纠正、同义词识别等各种自然语言处理函数,再构建用于高效查找的数据结构(例如倒排索引)。
在机器学习系统中,可以把模型看作是对训练数据应用各种特征提取和统计分析函数后得到的衍生物。模型应用于新的输入数据时,模型输出衍生自输入与模型(因此也间接衍生自训练数据)。
缓存通常以用户界面(UI)将要展示的形式保存数据聚合结果。因此,填充缓存需要知道 UI 引用了哪些字段;UI 的变更可能要求更新缓存填充方式的定义,并重建缓存。
二级索引的衍生函数需求极其常见,因此许多数据库都把它作为核心功能内置其中,只要执行 CREATE INDEX 就能调用。对于全文索引,常见语言的基础语言学功能或许内置于数据库,但更复杂的功能往往需要针对具体领域调整。在机器学习中,特征工程出了名地依赖具体应用,常常必须融入应用的用户交互与部署方式等详细知识 24。
如果创建衍生数据集的函数不像创建二级索引那样是标准化的套路,就需要用自定义代码处理应用特有的部分。而许多数据库恰恰难以应付这种自定义代码。关系数据库通常支持触发器、存储过程和用户定义函数,可以借此在数据库内部执行应用代码;但在数据库设计中,这些能力多少像是事后补上的。
分离应用代码与状态
理论上,数据库可以像操作系统一样,成为任意应用代码的部署环境。然而实践证明,数据库并不适合这个用途。依赖项与包管理、版本控制、滚动升级、可演化性、监控、指标、调用网络服务、与外部系统集成——数据库都无法很好地满足这些现代应用开发需求。
另一方面,Kubernetes、Docker、Mesos、YARN 等部署和集群管理工具,正是为运行应用代码而专门设计的。它们专注于做好这一件事,因此远胜于那些只把执行用户定义函数当作众多功能之一的数据库。
今天,大多数 Web 应用都以无状态服务的形式部署:任意用户请求都可以路由到任意应用服务器,服务器发出响应后便忘掉该请求的一切。这种部署方式很方便,因为服务器可以随意增加或移除;不过状态总得有个去处——通常是数据库。总体趋势是把无状态应用逻辑与状态管理(数据库)分开:既不把应用逻辑放进数据库,也不把持久状态放进应用 25。函数式编程社区喜欢开玩笑说:“我们信奉 教会(Church) 与 国家(state) 分离。”26
解释笑话通常会毁掉笑话,不过为了照顾没听懂的读者,这里还是解释一下。Church 指数学家阿隆佐・邱奇(Alonzo Church);他创造了 lambda 演算——一种早期计算形式,也是大多数函数式编程语言的基础。Lambda 演算没有可变状态(也就是没有可以被覆写的变量),所以也可以说,可变状态与 Church 的工作彼此分离。
在这种典型的 Web 应用模型中,数据库充当一种可以通过网络同步访问的可变共享变量。应用可以读取和更新这个变量,数据库则负责将它持久保存,并提供一定的并发控制与容错能力。
然而,在大多数编程语言中,你无法订阅可变变量的变化——只能定期读取它。与电子表格不同,变量的值发生变化时,读取者不会收到通知。(你可以在自己的代码中实现这种通知,这称为 观察者模式,observer pattern,但大多数语言都没有把这种模式作为内置功能。)
数据库继承了这种对待可变数据的被动方式:如果想知道数据库内容是否发生变化,你通常只能轮询(即定期重复查询)。订阅变更才刚刚开始成为数据库的一项功能。
数据流:状态变更与应用代码的相互作用
从数据流角度思考应用,意味着要重新协商应用代码与状态管理之间的关系。我们不再把数据库当作受应用操纵的被动变量,而是更加关注状态、状态变更以及处理这些变更的代码之间如何相互作用、彼此协作。应用代码响应一处的状态变更,并在另一处触发状态变更。
我们已经在变更数据捕获、Actor 模型、触发器和增量视图维护中见过这种想法。分拆数据库,就是把这一想法应用到主数据库之外的衍生数据集创建过程:缓存、全文检索索引、机器学习系统或分析系统。为此,我们可以使用流处理与消息传递系统。
维护衍生数据需要以下性质,而基于日志的消息代理能够提供这些性质:
维护衍生数据时,状态变更的顺序往往十分重要(如果多个视图都衍生自同一事件日志,它们就必须按相同顺序处理事件,才能彼此保持一致)。
容错必不可少:哪怕只丢失一条消息,衍生数据集也会与数据源永久失去同步。消息传递和衍生状态更新都必须可靠。
稳定的消息顺序和容错的消息处理要求相当严格,但它们的开销比分布式事务小得多,运维上也更加稳健。现代流处理器可以大规模提供这些顺序与可靠性保证,并允许应用代码作为流算子运行。
这类应用代码可以执行任意处理,补足数据库内置衍生函数通常不具备的能力。就像用管道串联起来的 Unix 工具一样,流算子也可以组合起来,围绕数据流构建大型系统。每个算子都以状态变更流作为输入,并产生其他状态变更流作为输出。
流处理器与服务
目前占主导地位的应用开发风格,是把功能拆分成一组 服务,服务之间通过 REST API 等同步网络请求通信。与单体应用相比,这种面向服务架构的主要优势,在于通过松散耦合实现组织上的可伸缩性:不同团队可以分别负责不同服务,从而减少团队间的协调工作(前提是各项服务能够独立部署和更新)。
把流算子组合成数据流系统,与微服务方法有许多相似之处 27、28。不过,两者底层的通信机制截然不同:前者使用单向异步消息流,后者使用同步的请求/响应交互。
除了 “事件驱动的架构” 中列出的优势(例如更好的容错性),数据流系统还可以获得优于传统 REST API 或 RPC 的性能。例如,假设顾客购买一件以一种货币定价、却用另一种货币付款的商品。要完成货币换算,就需要知道当前汇率。这个操作可以通过两种方式实现 27、29:
采用微服务方法时,处理购买操作的代码很可能查询汇率服务或数据库,以获得某种货币的当前汇率。
采用数据流方法时,处理购买操作的代码会预先订阅汇率更新流,并在汇率发生变化时把当前汇率记录到本地数据库中。真正处理购买操作时,只需查询本地数据库。
第二种方法用本地数据库查询,取代了对另一项服务的同步网络请求(本地数据库可能就在同一台机器上,甚至在同一个进程中)。在微服务方法中,也可以把汇率缓存在处理购买操作的服务本地,从而避免同步网络请求。然而,为了让缓存保持新鲜,就必须定期轮询汇率更新,或者订阅变更流——这恰好就是数据流方法所做的事。
数据流方法不仅更快,而且面对另一项服务发生故障时也更稳健。最快、最可靠的网络请求,就是根本不发出网络请求!RPC 不复存在,取而代之的是购买事件与汇率更新事件之间的流连接。
这种连接依赖于时间:如果日后重新处理购买事件,汇率早已改变。若想重建原始输出,就必须获得当初购买时的历史汇率。无论是查询服务还是订阅汇率更新流,都需要处理这种时间依赖性(参见 “连接的时间依赖性”)。
订阅变更流,而不是等到需要时才查询当前状态,使我们更接近电子表格式的计算模型:某项数据一旦变化,所有依赖它的衍生数据都能迅速更新。时间依赖连接等方面仍有许多悬而未决的问题,但围绕数据流理念构建应用,是一个很有前景、值得探索的方向。
观察衍生状态
从抽象层面看,上一节讨论的数据流系统提供了一套创建并持续更新衍生数据集(例如搜索索引、物化视图和预测模型)的过程。我们把这个过程称为 写路径(write path):每当有信息写入系统,它可能经过多轮批处理和流处理,最终所有衍生数据集都会更新,纳入这次写入的数据。图 13-1 展示了更新搜索索引的例子。

但你最初为什么要创建衍生数据集?很可能是为了日后查询。这就是 读路径(read path):处理用户请求时,从衍生数据集中读取数据,也许再对结果做些处理,最后构造返回给用户的响应。
写路径和读路径合在一起,涵盖了数据的完整旅程:从收集数据的地方,直至数据被消费的地方(很可能由另一个人消费)。写路径是预先计算的那一段——也就是数据一到达便立即完成,不管有没有人要求查看。读路径则只在有人请求时才发生。如果你熟悉函数式编程语言,或许会发现写路径类似于立即求值,读路径则类似于惰性求值。
如 图 13-1 所示,衍生数据集正是写路径与读路径相遇之处。它代表着写入时所需工作量与读取时所需工作量之间的权衡。
物化视图与缓存
全文检索索引就是一个很好的例子:写路径更新索引,读路径在索引中检索关键词。读写两端都要做些工作。写入需要更新文档中所有词项的索引条目;读取需要搜索查询中的每个词,并应用布尔逻辑,找出包含查询中 所有 词(AND 运算符)的文档,或者包含每个词的 任一 同义词(OR 运算符)的文档。
如果没有索引,搜索查询就必须扫描所有文档(类似 grep);文档数量一多,代价便会极其高昂。没有索引意味着写路径上的工作较少(无需更新索引),但读路径上的工作会多得多。
反过来,也可以设想预先计算所有可能查询的搜索结果。这样一来,读路径的工作就少了:无需进行布尔逻辑计算,只要找到相应查询的结果并返回即可。然而,写路径会昂贵得多:可能提出的搜索查询集合是无限的(或者至少随语料库中的词项数量呈指数增长),因此不可能预先计算所有搜索结果。
另一种选择是,只为一组固定的最常见查询预先计算搜索结果,使这些查询无需访问索引便能迅速得到响应;不常见的查询仍由索引处理。这通常称为常见查询的 缓存(cache),不过也可以称为物化视图:一旦出现应当纳入某项常见查询结果的新文档,它就必须随之更新。
这个例子说明,索引并不是写路径与读路径之间唯一可能的边界。既可以缓存常见搜索结果;文档数量较少时,也可以不用索引,进行类似 grep 的扫描。从这个角度看,缓存、索引与物化视图的作用很简单:它们移动了读路径与写路径之间的边界。我们通过预先计算结果,让写路径多做一些工作,从而节省读路径的开销。
写路径与读路径之间的工作边界,实际上正是 “案例研究:社交网络首页时间线” 中社交网络示例的主题。在那个例子中,我们还看到,名人与普通用户的读写路径边界可以划在不同位置。走过 500 页,我们又回到了原点!
有状态、可离线的客户端
写路径与读路径之间的边界很有意思,因为我们可以讨论如何移动这条边界,并探究这种移动在实践中意味着什么。下面换一个语境来看看这个想法。
过去,Web 浏览器是无状态客户端,只有接入互联网时才能做有用的事(离线时差不多只能在先前联网加载的页面里上下滚动)。然而,如今的单页 JavaScript Web 应用具备了许多有状态能力,包括客户端用户界面交互,以及 Web 浏览器内的持久化本地存储。移动应用同样可以在设备上保存大量状态,大多数用户交互也无需往返服务器。
我们在 “同步引擎与本地优先软件” 中看到,持久本地状态使一类应用成为可能:用户无需联网便可离线工作,有网络连接时再在后台与远程服务器同步 30。移动设备的蜂窝网络连接有时缓慢而不可靠;如果用户界面不用等待同步网络请求,而且应用大体上可以离线工作,对用户而言将是巨大优势。
当我们摆脱“无状态客户端与中央数据库通信”这一假设,转而在最终用户设备上维护状态时,一个充满新机会的世界便随之开启。尤其是,可以把设备上的状态视为 服务器状态的缓存。屏幕上的像素,是客户端应用中模型对象的物化视图;而模型对象,则是远程数据中心状态的本地副本 31。
将状态变更推送给客户端
在典型网页中,如果你用 Web 浏览器加载页面,随后服务器上的数据发生变化,那么除非重新加载页面,否则浏览器不会得知这一变化。浏览器只在某个时间点读取数据,并假定数据是静态的——它不会订阅服务器更新。因此,浏览器中的状态是一份陈旧缓存,除非明确轮询变更,否则不会更新。(RSS 等基于 HTTP 的 Feed 订阅协议,其实只是一种基本的轮询。)
较新的协议已经超越 HTTP 的基本请求/响应模式:服务器发送事件(EventSource API)与 WebSocket 提供了通信信道,让 Web 浏览器可以与服务器保持一条打开的 TCP 连接;只要连接仍在,服务器便能主动向浏览器推送消息。这样一来,服务器就能把客户端本地所存状态的任何变更主动告诉它,从而降低客户端状态的陈旧程度。
用写路径与读路径模型来说,主动把状态变更一路推送到客户端设备,意味着把写路径延伸至最终用户。客户端首次初始化时仍需要通过读路径取得初始状态,但此后便可依靠服务器发来的状态变更流。我们讨论过的流处理与消息传递理念,并不只限于在数据中心内运行:还可以继续向外延伸,直达最终用户设备 32。
设备有时会离线,在此期间无法收到服务器发出的任何状态变更通知。不过,这个问题我们已经解决过了:在 “消费者偏移量” 中,我们讨论了基于日志的消息代理的消费者如何在故障或断开后重新连接,并确保不漏掉断线期间到达的任何消息。同样的技术也适用于单个用户:每台设备都是一条小型事件流的订阅者。
端到端事件流
React 与 Elm 等用于开发有状态客户端和用户界面的工具 33,已经能够在底层状态发生变化时更新渲染出的用户界面。把这种编程模型进一步扩展,让服务器也能把状态变更事件推入客户端事件管道,是一件非常自然的事。
这样,状态变更就能沿端到端的写路径流动:从一台设备上触发状态变更的交互开始,经过事件日志、多个衍生数据系统与流处理器,一路到达另一台设备上观察该状态的用户界面。这些状态变更可以用相当低的延迟传播——例如端到端不到一秒。
一些应用(例如即时通讯与在线游戏)已经采用这种“实时”架构(这里指低延迟交互,而非响应时间保证)。但我们为什么不以这种方式构建所有应用?
挑战在于,无状态客户端与请求/响应交互的假设,已经深深嵌入我们的数据库、库、框架和协议。许多数据存储都支持一次请求返回一次响应的读写操作,但能够订阅变更的却少得多——也就是让一次请求随着时间推移返回一连串响应。
为了把写路径一直延伸到最终用户,我们必须从根本上重新思考许多系统的构建方式:离开请求/响应交互,转向发布/订阅数据流 31。这需要付出努力,但也会带来响应更灵敏的用户界面,以及更好的离线支持。
读也是事件
前面说过,当流处理器把衍生数据写入某个存储(数据库、缓存或索引),而这个存储随后接受查询时,它就充当了写路径与读路径之间的边界。该存储允许对数据进行随机访问读取;否则,读取查询就必须扫描整条事件日志。
在许多情况下,数据存储与流处理系统彼此分离。但别忘了,流处理器也需要维护状态,才能执行聚合与连接。这种状态通常隐藏在流处理器内部,不过有些框架也允许外部客户端查询它 34,从而让流处理器本身变成一种简单的数据库。
我们再把这个想法推进一步。到目前为止,写入通过事件日志进入存储,而读取则是短暂的网络请求,直接发往保存待查数据的节点。这种设计很合理,却不是唯一选择。我们也可以把读取请求表示成事件流,把读事件和写事件都送入流处理器;处理器通过向输出流发出读取结果,来响应读事件 35。
当读写都表示为事件,并路由到同一个流算子处理时,我们实际上是在读取查询流与数据库之间执行流表连接。读事件需要发送到保存相应数据的数据库分片,就像批处理器和流处理器执行连接时,需要按同一个键对输入进行协同分区一样。
处理请求与执行连接之间的这种对应关系,是一个非常基础的概念 36。一次性读取请求穿过连接算子后,算子立刻将它忘掉;订阅请求则是一项持久连接,与连接另一侧过去和未来的事件不断匹配。
记录读事件日志,在追踪系统中的因果依赖与数据溯源方面也可能有益:它让你能够重建用户在作出某项决定前看到了什么。例如,在网上商店中,向顾客显示的预计送达日期与库存状态,很可能会影响他们是否选择购买一件商品 4。要分析这种关联,就必须记录用户对送货与库存状态查询的结果。
因此,把读取请求写入持久存储,有助于更好地追踪因果依赖,但会增加存储与 I/O 成本。如何优化这类系统、降低开销,仍是一个开放的研究问题 2。不过,如果你本来就出于运维目的,在处理请求时顺带记录了读取请求日志,那么反过来让日志成为请求来源,并不算多么巨大的改变。
多分片数据处理
对于只涉及单个分片的查询,通过流发送查询并收集响应或许有些小题大做。不过,这个想法开启了一种可能:利用流处理器已经具备的消息路由、分片与连接基础设施,分布式执行需要组合多个分片数据的复杂查询。
Storm 的分布式 RPC 功能支持这种使用模式。例如,它曾被用于计算社交网络上看过某个 URL 的人数——也就是发布过该 URL 的所有用户,其关注者集合的并集 37。由于用户集合经过分片,这项计算必须组合许多分片的结果。
这种模式的另一个例子是欺诈防范:为了评估某项购买事件是否可能存在欺诈,可以检查用户 IP 地址、电子邮件地址、账单地址、送货地址等各自的信誉评分。每个信誉数据库本身都经过分片,因此,为某项购买事件收集这些评分,需要依次连接多个采用不同分片方式的数据集 38。
数据仓库查询引擎内部的查询执行图,也具有类似特征。如果需要执行这种多分片连接,使用原生提供该功能的数据库,很可能比借助流处理器自行实现更简单。不过,把查询视为流,仍为构建逼近传统现成方案能力极限的大规模应用提供了一种选择。
追求正确性
对于只读取数据的无状态服务,出了问题也没什么大不了:修复缺陷、重启服务,一切便恢复正常。数据库等有状态系统却没这么简单:它们的设计目标是(近乎)永久保存信息,因此一旦出了问题,影响也可能永远持续——这意味着我们必须更加仔细地思考 39。
我们希望构建既可靠又 正确 的应用(也就是即使面对各种故障,程序的语义仍有明确的定义,也能为人理解)。大约四十年来,原子性、隔离性和持久性等事务属性,一直是构建正确应用的首选工具。然而,这套基础并不像看起来那样牢固——弱隔离级别所引发的困惑便是一例(参见 “弱隔离级别”)。
在某些领域,事务已被彻底抛弃,取而代之的是性能与可伸缩性更好、语义却混乱得多的模型。人们经常谈论 一致性,却很少把它定义清楚。有些人声称,为了更高的可用性,我们应该“拥抱弱一致性”,却说不清这在实践中究竟意味着什么。
对于如此重要的主题,我们的理解与工程方法却出奇地不可靠。例如,要判断某个应用在特定事务隔离级别或复制配置下运行是否安全,非常困难 40、41。简单方案在并发度低、没有故障时往往看似正确,一旦环境要求更高,便会暴露出许多隐蔽缺陷。
例如,Kyle Kingsbury 的 Jepsen 实验 42 揭示了某些产品宣称的安全保证,与它们遭遇网络问题和崩溃时的实际行为之间存在巨大差距。即便数据库等基础设施产品本身毫无问题,应用代码仍须正确使用它们提供的功能;如果配置本就难以理解(弱隔离级别、法定人数配置等都是如此),这个过程很容易出错。
如果你的应用能够容忍偶尔以不可预测的方式损坏或丢失数据,事情会简单得多;或许只要祈求好运,就能勉强应付。反过来,如果你需要更强的正确性保证,可串行化与原子提交是成熟的方法,却也代价不菲:它们通常只能在单个数据中心内工作(排除了地理分布式架构),还会限制系统能够达到的规模与容错能力。
传统事务方法并不会消失,但要让应用既正确、又能抵御故障,事务并不是最终答案。本节将探讨在数据流架构语境下思考正确性的几种方式。
数据库的端到端原则
应用使用了具有较强安全属性的数据系统(例如可串行化事务),并不等于它就不会丢失或损坏数据。例如,如果应用缺陷导致它写入错误数据,或者从数据库删除数据,可串行化事务也救不了你。这正是支持不可变与仅追加数据的一条理由:如果不让错误代码拥有摧毁正确数据的能力,从这类错误中恢复就容易得多。
尽管不可变性很有用,但它本身并非万灵药。下面来看一个更隐蔽的数据损坏例子。
恰好一次执行操作
在 “容错” 中,我们见过 恰好一次(或 等效一次)语义。如果处理消息时出了问题,可以选择放弃(丢弃消息,也就是造成数据丢失),也可以重试。选择重试,就要承担一种风险:第一次其实已经成功,只是你没能得知,于是消息最终被处理两次。
处理两次也是一种数据损坏:我们不希望同一项服务向顾客收取两次费用(多收费),也不希望计数器递增两次(夸大指标)。在这里,恰好一次 是指这样安排计算:即便某项操作确实由于故障而重试,最终效果也与从未发生故障时相同。前面已经讨论过几种实现方式。
最有效的方法之一,是让操作变得 幂等:也就是无论执行一次还是多次,效果都相同。然而,要把本来并不幂等的操作改造成幂等操作,需要额外投入并谨慎处理:你可能需要维护额外的元数据(例如曾经更新过某个值的操作 ID 集合),并确保从一个节点故障切换到另一个节点时使用栅栏机制(参见 “分布式锁和租约”)。
抑制重复
这种需要抑制重复的模式,还会出现在流处理之外的许多地方。例如,TCP 使用数据包的序列号,让接收方按正确顺序排列数据包,并判断网络中是否有包丢失或重复。丢失的包会被重传,重复的包则由 TCP 栈移除,之后数据才交给应用。
不过,这种重复抑制只能在单条 TCP 连接的范围内发挥作用。假设这条 TCP 连接把客户端连到数据库,当前正在执行 例 13-1 中的事务。在许多数据库中,事务与客户端连接绑定(如果客户端发送多条查询,数据库之所以知道它们属于同一事务,是因为它们都从同一条 TCP 连接发来)。如果客户端发出 COMMIT 之后、尚未收到数据库服务器的回应之前遭遇网络中断与连接超时,它便无从知道事务究竟已经提交还是已经中止(图 9-1)。
例 13-1. 从一个账户向另一个账户非幂等地转账
客户端可以重新连接数据库并重试事务,但这已经超出 TCP 重复抑制的范围。由于 例 13-1 中的事务并不幂等,结果可能实际转了 $22,而不是预期的 $11。因此,尽管 例 13-1 是说明事务原子性的标准示例,它实际上并不正确,真正的银行也不会这样工作 3。
两阶段提交协议(参见 “两阶段提交(2PC)”)打破了 TCP 连接与事务之间的一一对应,因为它必须允许事务协调者在网络故障后重新连接数据库,并告诉数据库应当提交还是中止悬而未决的事务。这足以保证事务只执行一次吗?很遗憾,并不足够。
即使能够抑制数据库客户端与服务器之间的重复事务,我们仍然要担心最终用户设备与应用服务器之间的网络。例如,如果最终用户客户端是 Web 浏览器,它可能通过 HTTP POST 请求向服务器提交指令。也许用户的蜂窝数据连接很差:POST 成功发出,但信号在服务器响应到达之前变得太弱。
此时,用户很可能看到错误消息,并手动重试。Web 浏览器会警告:“确定要再次提交此表单吗?”——用户选择确定,因为他们本来就希望操作发生。(Post/Redirect/Get 模式 43 可以在正常情况下避免这条警告,却无法解决 POST 请求超时。)在 Web 服务器看来,重试是另一次请求;在数据库看来,它也是另一个事务。通常的去重机制无济于事。
唯一标识请求
要让请求经过多跳网络通信后仍然幂等,只依靠数据库提供的事务机制是不够的——必须考虑请求的 端到端 流程。
例如,可以为请求生成唯一标识符(如 UUID),将它作为隐藏表单字段放入客户端应用;也可以对所有相关表单字段计算哈希,衍生出请求 ID 3。如果 Web 浏览器提交两次 POST 请求,两次请求就会携带同一个请求 ID。随后,可以把这个请求 ID 一路传到数据库,并检查给定 ID 的请求永远只执行一次,如 例 13-2 所示。
例 13-2. 使用唯一 ID 抑制重复请求
例 13-2 依赖 request_id 列上的唯一性约束。如果事务试图插入一个已经存在的 ID,INSERT 就会失败,事务随之中止,避免再次生效。即便隔离级别较弱,关系数据库通常也能正确维护唯一性约束(而应用层面的“先检查再插入”在非可串行化隔离下可能失效,如 “写偏差与幻读” 所述)。
除了抑制重复请求,例 13-2 中的 requests 表还充当了一种事件日志,可用于事件溯源或变更数据捕获。账户余额的更新其实不必与插入事件发生在同一事务中,因为它们是冗余状态,可以由下游消费者从请求事件衍生出来——只要事件得到恰好一次处理,而这一点同样可以借助请求 ID 强制保证。
端到端原则
抑制重复事务的情形,只是一个更普遍原则的例子。这个原则称为 端到端原则(end-to-end argument),由 Saltzer、Reed 和 Clark 于 1984 年提出 44:
只有借助位于通信系统两端的应用所掌握的知识与提供的协助,所讨论的功能才能得到完整、正确的实现。因此,不可能把这一功能作为通信系统本身的一项功能来提供。(有时,通信系统提供的不完整版本可以用来提升性能。)
在我们的例子中,所讨论的功能 是重复抑制。TCP 会在 TCP 连接层面抑制重复数据包,一些流处理器则在消息处理层面提供所谓的恰好一次语义;但如果第一次请求超时,这些机制都不足以防止用户再次提交同一请求。TCP、数据库事务和流处理器自身都无法彻底排除这种重复。解决问题需要端到端方案:让事务标识符从最终用户客户端一路传到数据库。
端到端原则同样适用于检查数据完整性:以太网、TCP 与 TLS 内置的校验和可以检测网络数据包损坏,却无法检测网络连接两端收发软件中的缺陷所造成的损坏,也无法检测存储数据的磁盘发生的损坏。如果想捕捉所有可能的数据损坏来源,就还需要端到端校验和。
类似的道理也适用于加密 44:家庭 WiFi 网络的密码可以阻止他人窃听你的 WiFi 流量,却防不住互联网其他地方的攻击者;客户端与服务器之间的 TLS/SSL 可以抵御网络攻击者,却防不住服务器本身遭到入侵。只有端到端加密与认证才能抵御所有这些威胁。
尽管低层功能(TCP 重复抑制、以太网校验和、WiFi 加密)单凭自身无法提供所需的端到端功能,它们仍然有用,因为它们降低了高层发生问题的概率。例如,如果没有 TCP 把数据包重新排成正确顺序,HTTP 请求往往会变得支离破碎。我们只需记住:仅靠低层可靠性功能,并不足以保证端到端正确性。
在数据系统中应用端到端思维
这又把我们带回最初的论点:应用使用了提供较强安全属性的数据系统(例如可串行化事务),并不代表它一定不会丢失或损坏数据。应用自身同样必须采取端到端措施,例如抑制重复。
这未免令人遗憾,因为容错机制很难正确实现。TCP 等低层可靠性机制相当有效,因此剩余的高层故障很少发生。若能把余下的高层容错机制封装成一种抽象,让应用代码不必操心,那再好不过——但我们似乎还没有找到合适的抽象。
长期以来,事务一直被视为有用的抽象。正如 第 8 章 所述,它把各种可能的问题(并发写入、违反约束、崩溃、网络中断、磁盘故障)压缩成两种可能结果:提交或中止。这极大简化了编程模型,但仍然不够。
事务的代价很高,涉及异构存储技术时尤其如此(参见 “跨不同系统的分布式事务”)。当我们因为分布式事务代价过高而拒绝使用它时,最终不得不在应用代码中重新实现容错机制。本书大量例子表明,对并发与部分失效进行推理既困难又反直觉,因此大多数应用级机制都无法正确工作,结果便是数据丢失或损坏。
因此,值得探索这样的容错抽象:既能轻松提供应用特有的端到端正确性,又能在大规模分布式环境中保持良好的性能与运维特性。
强制约束
下面结合数据库分拆的理念来思考正确性。前面看到,只要把请求 ID 从客户端一路传到记录写入的数据库,就能实现端到端重复抑制。那么,其他类型的约束呢?
我们特别关注 唯一性约束——也就是 例 13-2 所依赖的约束。在 “约束与唯一性保证” 中,我们还见过另外几种需要强制唯一性的应用功能:用户名或电子邮件地址必须唯一标识一名用户;文件存储服务不能有多个同名文件;两个人不能预订同一个航班座位或剧院座位。
其他约束也十分相似,例如确保账户余额永不为负、售出的商品不超过仓库库存,或者会议室的预订时间不能重叠。强制唯一性的技术通常也能用于这类约束。
唯一性约束需要共识
我们在 第 10 章 中看到,在分布式环境中强制唯一性约束需要共识:如果有多个取值相同的并发请求,系统必须设法决定接受哪一个冲突操作,并以违反约束为由拒绝其余操作。
达成这种共识最常见的方法,是让单个节点担任领导者,负责作出所有决定。只要你不介意让全部请求汇入单个节点(哪怕客户端身处地球另一端),而且该节点不发生故障,这种方法就行得通。Raft 等共识算法则解决了当前领导者失效(或者因网络问题而被认为失效)时,如何安全选举新领导者并避免脑裂的问题。
唯一性检查可以根据必须唯一的值进行分片,从而横向扩展。例如,如果要像 例 13-2 那样按请求 ID 保证唯一性,就可以确保所有具有相同请求 ID 的请求都路由到同一个分片;如果用户名必须唯一,则可以按用户名的哈希值分片。
不过,异步多主复制并不适用,因为不同领导者可能并发接受相互冲突的写入,使值不再唯一。如果希望立即拒绝任何违反约束的写入,同步协调便不可避免 45。
基于日志的消息传递中的唯一性
共享日志保证所有消费者以相同顺序看到消息——这种保证在形式上称为 全序广播,并且等价于共识(参见 “共识的多面性”)。在使用基于日志消息传递的数据库分拆方法中,可以采用非常相似的方式强制唯一性约束。
流处理器用单个线程依次消费某个日志分片中的所有消息。因此,如果日志按照必须唯一的值来分片,流处理器就能毫无歧义地、确定性地判断,多项冲突操作中哪一项最先出现在日志里。例如,多个用户试图注册同一个用户名时 46:
每个用户名请求都被编码成一条消息,并追加到由用户名哈希值决定的分片。
流处理器依次读取日志中的请求,用本地数据库记录哪些用户名已经被占用。每当有请求申请一个可用用户名,处理器就把该名称记为已占用,并向输出流发出成功消息;每当请求申请一个已经占用的用户名,处理器就向输出流发出拒绝消息。
请求用户名的客户端观察输出流,等待与自身请求对应的成功或拒绝消息。
这一算法与我们在 第 10 章 中见过的、使用共享日志实现共识的构造相同。只要增加分片数量,就能轻松扩展到很高的请求吞吐量,因为每个分片都可以独立处理。
这种方法不仅适用于唯一性约束,也适用于许多其他约束。其基本原则是:任何可能冲突的写入都路由到同一个分片,依次处理。冲突的定义可能取决于具体应用,但流处理器可以用任意逻辑来验证请求。
多分片请求处理
当一项操作涉及多个分片时,要让它在满足约束的同时原子执行,问题就更有意思了。例 13-2 可能涉及三个分片:保存请求 ID 的分片、保存收款方账户的分片,以及保存付款方账户的分片。这三者彼此独立,没有理由一定落在同一个分片中。
采用传统数据库方法时,执行这项事务需要跨三个分片原子提交;这实质上迫使它与这三个分片中任何一个分片上的其他所有事务形成全序。既然出现了跨分片协调,各分片便无法再独立处理,吞吐量很可能因此受损。
然而,借助分片日志与流处理器,无需跨分片事务也能实现等价的正确性。图 13-2 展示了一项付款事务:先检查源账户余额是否充足;若余额充足,则在扣除手续费的同时,把一笔金额原子地转入目标账户。具体过程如下 47:

用户客户端为从源账户向目标账户转账的请求赋予唯一请求 ID,再根据源账户 ID,把请求追加到相应日志分片。
流处理器读取请求日志,并维护一个数据库,其中保存源账户的状态以及已经处理过的请求 ID。这个数据库的内容完全衍生自日志。当流处理器遇到一个从未见过的请求 ID 时,它会在本地数据库中检查源账户余额是否足以完成转账。
如果余额充足,处理器就在本地数据库中更新源账户状态,预留付款金额,并向另外几条日志发出事件:向源账户的日志分片(也就是处理器自己的输入日志)发出一条出账事件,向目标账户的日志分片发出一条入账事件,再向手续费账户的日志分片发出一条入账事件。发出的这些事件都包含原始请求 ID。
出账事件最终会回到源账户处理器(其间可能已经收到一些不相干的事件)。流处理器根据请求 ID 认出,这是先前已经预留的一笔付款,于是现在执行付款,再次更新本地保存的源账户状态。它会根据请求 ID 忽略重复事件。
目标账户与手续费账户的日志分片分别由独立的流处理任务消费。它们收到入账事件后,会更新各自的本地状态以反映这笔款项,并根据请求 ID 对事件去重。
图 13-2 把三个账户画在三个不同分片中,但它们也完全可以位于同一个分片——这无关紧要。我们只需要保证:给定账户的所有事件都严格按照日志顺序处理,并采用至少一次语义,而且流处理器是确定性的。
例如,设想源账户处理器在处理付款请求时崩溃。崩溃之前,输出消息可能已经发出,也可能尚未发出。处理器从崩溃中恢复后,会再次处理同一个请求(因为采用至少一次语义);由于处理器具有确定性,它仍会对是否允许付款作出相同决定。因此,它会向出账、入账和手续费账户分片发出带有相同请求 ID 的同一批输出消息。如果这些消息是重复的,下游消费者会根据请求 ID 将其忽略。
这个系统的原子性不来自任何事务,而来自把初始请求事件写入源账户日志这一原子操作。一旦那条事件进入日志,所有下游事件最终也都会写入——可能要等流处理器从崩溃中恢复,可能还会出现重复,但它们终究会出现。
如果采用恰好一次语义,这个例子实现起来会更容易,因为它能确保流处理器的本地状态与已经处理的消息集合保持一致。因此,如果处理器崩溃并重新处理某些消息,它的本地状态也会重置到处理这些消息之前的状态。
如果 图 13-2 中的用户想知道转账是否获批,可以订阅源账户的日志分片,等待出账事件。如果希望在余额不足时明确通知用户,流处理器可以向该日志分片发出一条“付款被拒”事件。
通过把多分片事务拆成多个采用不同分片方式的阶段,并使用端到端请求 ID,我们实现了相同的正确性属性(每个请求对付款方与收款方账户都恰好应用一次):即使发生故障也不例外,而且无需使用原子提交协议。
及时性与完整性
许多事务系统都有一项便利的性质:一个事务提交后,其写入立刻对其他事务可见。这项性质的形式化名称是 严格可串行化(strict serializability;参见 “线性一致性与可串行化”)。
把一项操作分拆为流处理器的多个阶段后,情况却并非如此:日志消费者在设计上就是异步的,所以发送者不会等待消费者处理完自己的消息。不过,客户端仍可以等待某条消息出现在输出流上。例如,图 13-2 中的用户可以等待出账事件或付款被拒事件,这取决于源账户中是否有足够资金。
在这个例子中,检查源账户余额是否正确,并不取决于发出请求的用户是否等待结果。等待只是为了同步告知用户付款是否成功,这项通知与处理请求产生的效果彼此解耦。
更一般地说,一致性(consistency)这个术语混合了两种不同的需求,而它们值得分开考虑:
- 及时性(timeliness)
及时性是指确保用户观察到系统的最新状态。前面看到,如果用户从陈旧的数据副本中读取,可能观察到不一致的系统状态(参见 “复制延迟的问题”)。不过,这种不一致只是暂时的,只需等待并重试,最终便会消失。
CAP 定理中的一致性指线性一致性,这是实现及时性的一种强保证。写后读一致性(read-after-write consistency)等较弱的及时性属性同样有用。
- 完整性(integrity)
完整性是指没有损坏:既不丢失数据,也没有相互矛盾或虚假的数据。尤其是,如果某个衍生数据集作为底层数据之上的视图来维护,衍生过程必须正确。例如,数据库索引必须准确反映数据库内容——漏掉某些记录的索引没什么用。
如果完整性遭到破坏,不一致就是永久的:大多数情况下,等待和重试无法修复数据库损坏,必须明确进行检查和修复。在 ACID 事务语境中,“一致性”通常被理解为某种应用特有的完整性概念。原子性与持久性是维护完整性的重要工具。
用一句口号来说:违反及时性叫“最终一致”,违反完整性则叫“永远不一致”。
在大多数应用中,完整性都比及时性重要得多。及时性遭到破坏会令人烦恼和困惑,完整性遭到破坏却可能带来灾难。
例如,信用卡账单上没有出现过去 24 小时内完成的一笔交易,并不会令人意外——这类系统存在一定延迟很正常。我们知道银行会异步对账和结算交易,因此这里的及时性并不重要 3。然而,如果账单余额不等于交易总额加上上期账单余额(求和出错),或者一笔交易向你扣了款、商户却没有收到钱(资金凭空消失),问题就非常严重。这些都是对系统完整性的破坏。
数据流系统的正确性
ACID 事务通常同时提供及时性保证(例如线性一致性)和完整性保证(例如原子提交)。因此,如果从 ACID 事务的角度看待应用正确性,及时性与完整性之间的区别便无关紧要。
另一方面,本章讨论的基于事件的数据流系统有一项有趣性质:它们把及时性与完整性解耦了。异步处理事件流时,除非明确构建消费者,让它等到消息到达后才返回,否则就没有及时性保证。例如,用户可以请求一笔付款,随后在流处理器执行该请求之前读取自己的账户状态;此时,用户看不到刚刚请求的付款。
然而,完整性事实上是流式系统的核心。恰好一次(exactly-once)或 等效一次(effectively-once)语义就是维护完整性的一种机制。事件丢失或生效两次,都可能破坏数据系统的完整性。因此,面对故障时,容错消息传递与重复抑制(例如幂等操作)是维护数据系统完整性的关键。
正如上一节所见,可靠的流处理系统无需分布式事务与原子提交协议也能保持完整性。这意味着它们有望实现同等程度的正确性,同时获得好得多的性能与运维稳健性。我们通过组合以下机制实现了这种完整性:
把写入操作的内容表示成一条消息,使其能够轻松原子写入——这种方式非常适合事件溯源
使用确定性衍生函数,从这一条消息衍生出所有其他状态更新,类似于存储过程
让客户端生成的请求 ID 贯穿所有处理层次,从而实现端到端重复抑制与幂等性
让消息不可变,并允许不时重新处理衍生数据,从而更容易从缺陷中恢复
宽松解释约束
如前所述,强制唯一性约束需要共识,通常通过让某个分片中的所有事件都汇入单个节点来实现。如果希望采用传统形式的唯一性约束,这项限制便无法避免,流处理也绕不过去。
不过还要意识到:在许多真实应用中,业务需求其实允许违反那些看似硬性约束的规则:
如果顾客订购的商品超过仓库库存,可以补订货物,为延误向顾客道歉,并提供折扣。其实,即便只是叉车碾坏了仓库里的部分商品,导致实际库存少于预期,你也必须这样处理 3。因此,为了应对叉车事故,道歉工作流本来就必须纳入业务流程;对库存数量设置硬约束或许没有必要。
同样,许多航空公司会超卖机票,因为预计有些乘客会误机;许多酒店也会超卖客房,因为预计有些客人会取消预订。在这些情形中,企业出于业务考虑故意违反“一座一人”的约束,并设置补偿流程(退款、升级、在附近酒店免费提供房间),处理需求超出供给的情况。即便没有超卖,也需要道歉与补偿流程,以应对恶劣天气或员工罢工造成的航班取消——从这类问题中恢复本来就是正常业务的一部分 3。
如果有人取出的金额超过账户余额,银行可以收取透支费,并要求对方偿还欠款。只要限制每日提款总额,银行承担的风险就有上限。
在跨组织集成数据的系统中,不一致不可避免,因此必须有修正机制来处理它们。正如 “批处理用例” 所指出的,银行之间的付款结算就是一个例子。
因此,在许多业务场景中,暂时违反约束、稍后再通过道歉修正,是可以接受的。这种用于纠正错误的变更称为 补偿性事务(compensating transaction)48、49。道歉的代价各不相同(金钱或声誉上的代价),却往往很低:已经发出的电子邮件无法撤回,但可以再发一封邮件更正;信用卡不慎扣款两次,可以退回其中一笔,代价只是手续费,或许再加上一项顾客投诉。ATM 一旦吐出现金,确实无法直接收回;但原则上,如果账户已经透支而顾客拒绝还款,可以派催收人员追回欠款。
道歉的代价能否接受,是一项业务决策。如果能够接受,那么“写入数据前先检查全部约束”的传统模型就限制过多。完全可以先乐观地执行写入,再事后检查约束。对于那些一旦出错便很难挽回的事情,仍可以确保在执行前完成验证;但这并不意味着,连数据写入之前也必须先做验证。
这些应用 确实 需要完整性:谁也不希望预订记录丢失,或因借贷不匹配而让资金凭空消失。但是,它们在强制约束时 并不需要 及时性:如果售出的商品超过仓库库存,事后道歉并补救即可。这与我们在 “处理写入冲突” 中讨论的冲突解决方法相似。
避免协调的数据系统
我们现在已经做了两个有趣的观察:
数据流系统无需原子提交、线性一致性或同步的跨分片协调,也能维护衍生数据的完整性保证。
尽管严格的唯一性约束需要及时性与协调,许多应用其实可以接受宽松约束:只要完整性始终得到维护,约束可以暂时遭到违反,稍后再修复。
把这两点结合起来就意味着:数据流系统无需协调,便可为许多应用提供数据管理服务,同时仍给出强有力的完整性保证。这种 避免协调(coordination-avoiding)的数据系统极具吸引力:与需要同步协调的系统相比,它们能获得更好的性能与容错能力 45。
例如,这类系统可以采用多主配置,分布在多个数据中心,并在区域之间异步复制。任何一个数据中心都能独立于其他数据中心继续运行,因为不需要跨区域同步协调。这样的系统只提供较弱的及时性保证——不引入协调就不可能实现线性一致性——却仍能提供强有力的完整性保证。
在这种情况下,可串行化事务作为维护衍生状态的一环仍然有用,不过可以把它限制在自己擅长的小范围内 6。不再需要 XA 等异构分布式事务。仍然可以在确有需要之处引入同步协调(例如在执行无法恢复的操作之前,强制实施严格约束);但如果应用中只有一小部分需要协调,就没必要让所有部分都付出协调成本 32。
也可以换个角度看协调与约束:它们减少了因不一致而道歉的次数,却也可能降低系统性能和可用性,从而增加因服务中断而道歉的次数。道歉次数不可能降到零,但可以根据自身需要寻找最佳平衡点——既不会出现太多不一致,也不会遇到太多可用性问题。
信任但验证
前面关于正确性、完整性与容错的所有讨论,都建立在一组假设之上:某些事情可能出错,另一些事情不会。我们把这些假设称为 系统模型(system model;参见 “系统模型与现实”)。例如,我们应当假设进程可能崩溃、机器可能突然断电、网络可能任意延迟或丢弃消息;但也可能假设,写入磁盘的数据经过 fsync 后不会丢失、内存中的数据不会损坏、CPU 的乘法指令总能返回正确结果。
这些假设相当合理,因为绝大多数时候它们都成立;如果必须时刻担心计算机会算错,我们将寸步难行。传统系统模型以二元方式看待故障:假设有些事情可能发生,另一些事情绝不可能发生。现实却更像是概率问题:有些事情更常见,有些事情较少见。真正的问题是,违反假设的情况是否频繁到我们会在实践中遇见。
我们已经看到,数据可能在内存中损坏(参见 “硬件与软件故障”),可能在磁盘上损坏(参见 “复制与持久性”),也可能在网络中损坏(参见 “弱形式的撒谎”)。或许我们应该更加重视这一点?当系统规模足够大,再小概率的事情也会发生。
面对软件缺陷时维护完整性
除了这类硬件问题,软件缺陷也始终是一项风险,而低层的网络、内存或文件系统校验和捕捉不到它们。即便广泛使用的数据库软件也存在缺陷:例如,过去某些版本的 MySQL 未能正确维护唯一性约束 50,PostgreSQL 的可串行化隔离级别过去也曾出现写偏差异常 51。MySQL 与 PostgreSQL 都是稳健且口碑良好的数据库,经过许多人多年的实战检验;不够成熟的软件,情况很可能糟得多。
尽管人们投入大量精力仔细设计、测试和审查,缺陷仍会悄然混入。它们虽然罕见,最终也会被发现和修复,却仍然存在一段可能损坏数据的窗口期。
至于应用代码,我们必须假设其中有更多缺陷,因为大多数应用接受的审查与测试,远远不及数据库代码。许多应用甚至没有正确使用数据库提供的完整性维护功能,例如外键或唯一性约束 25。
ACID 意义上的一致性,建立在这样一种想法之上:数据库从一致状态开始,事务再把它从一个一致状态转变为另一个一致状态。因此,我们期望数据库始终处于一致状态。然而,只有假设事务没有缺陷,这种说法才有意义。如果应用以某种错误方式使用数据库——例如不安全地采用弱隔离级别——数据库的完整性便无法保证。
不要盲信承诺
硬件和软件都不总能达到理想状态,因此数据损坏迟早似乎不可避免。至少,我们应该有办法发现数据已经损坏,从而修复它,并努力追查错误来源。检查数据完整性的过程称为 审计(auditing)。
正如 “不可变事件的优点” 所述,审计并不只适用于财务应用。不过,可审计性在金融领域格外重要,恰恰因为人人都知道错误难免发生,也都认可能够发现并修复问题的必要性。
成熟系统同样倾向于考虑小概率故障的可能性,并管理这种风险。例如,HDFS 和 Amazon S3 等大规模存储系统不会完全信任磁盘:它们运行后台进程,不断回读文件、与其他副本比较,并把文件从一个磁盘移动到另一个磁盘,以降低静默损坏的风险 52、53。
如果想确认数据仍然存在,就必须真正读取并检查。绝大多数时候数据依然完好;但万一不是,你肯定希望越早发现越好。同理,不时尝试从备份恢复也很重要——否则,你可能直到数据已经丢失、为时已晚,才发现备份根本无法使用。不要盲信一切都在正常工作。
HDFS 与 S3 仍然必须假设磁盘在绝大多数时候能够正确工作——这个假设很合理,却不同于假设磁盘 始终 正确工作。然而,目前采用这种“信任,但要验证”方式持续自我审计的系统并不多。许多系统假定正确性保证是绝对的,完全没有为罕见的数据损坏预作安排。未来,我们或许会看到更多 自我验证(self-validating)或 自我审计(self-auditing)系统:它们不断检查自身完整性,而不是依赖盲目信任 54。
为可审计性而设计
如果一个事务修改了数据库中的多个对象,事后很难看出这项事务究竟意味着什么。即便捕获了事务日志,各个表里的插入、更新与删除也未必能清楚说明,为什么 要执行这些修改。当初决定作出这些修改的应用逻辑调用转瞬即逝,无法重现。
相比之下,基于事件的系统可以提供更好的可审计性。在事件溯源方法中,系统的用户输入被表示为一条不可变事件,由此产生的所有状态更新都衍生自这条事件。衍生过程可以做到确定性与可重复性,因此,用同一版本的衍生代码处理同一份事件日志,便会得到相同的状态更新。
明确表示数据流,能让 数据溯源(data provenance)清晰得多,从而使完整性检查更切实可行。对于事件日志,可以用哈希检查事件存储是否遭到损坏;对于任何衍生状态,可以重新运行当初从事件日志衍生它的批处理器与流处理器,检查是否得到同样结果,甚至还可以并行运行一条冗余的衍生流程。
具有确定性且定义明确的数据流,也让系统执行过程更容易调试和追踪,从而查明系统 为什么 做了某件事 4、55。如果发生意外,能够重现导致意外事件的确切情境将极有价值——这是一种时间旅行式调试能力。
再谈端到端原则
如果无法完全相信系统中的每个组件都不会造成损坏——每一件硬件都不出故障,每一段软件都没有缺陷——那么至少必须定期检查数据完整性。如果不检查,往往要等损坏造成下游影响、为时已晚时才会发现;那时追查问题将困难得多,代价也高得多。
检查数据系统完整性,最好采用端到端方式:完整性检查涵盖的系统越多,处理流程某个阶段的损坏就越不容易逃过检查。如果能够检查整条衍生数据管道端到端的正确性,那么路径上的所有磁盘、网络、服务与算法,也都隐含在检查范围之内。
持续进行端到端完整性检查,会增强你对系统正确性的信心,从而让你行动得更快 56。审计与自动化测试一样,提高了迅速发现缺陷的概率,因而降低系统变更或采用新存储技术造成损害的风险。如果不再害怕作出改变,就能更好地推动应用演化,以满足不断变化的需求。
可审计数据系统的工具
目前,很少有数据系统把可审计性当作首要目标。一些应用会实现自己的审计机制,例如把所有变更记录到独立的审计表;然而,要保证审计日志与数据库状态的完整性仍然很难。可以借助硬件安全模块定期签名,让事务日志不可篡改,但这并不能保证一开始进入日志的就是正确事务。
Bitcoin 或 Ethereum 等区块链,是带有密码学一致性检查的共享仅追加日志;其中存储的交易是事件,智能合约基本上就是流处理器。它们使用的共识协议确保所有节点对同一事件序列达成一致。与 第 10 章 的共识协议不同,区块链具有拜占庭容错能力:即便某些参与节点的数据已经损坏,系统仍能继续工作,因为各个副本会不断互相检查完整性。
对于大多数应用,区块链的开销太高,并不实用。不过,其中一些密码学工具也能用于更轻量的环境。例如,默克尔树 57 是由哈希构成的树,可以高效证明某条记录出现在某个数据集中(也能证明其他一些事情)。证书透明性 使用经过密码学验证的仅追加日志与默克尔树,检查 TLS/SSL 证书的有效性 58、59;它让每条日志由单个领导者负责,从而无需共识协议。
未来,证书透明性与分布式账本所采用的完整性检查和审计算法,或许会在一般数据系统中得到更广泛的应用。要让它们拥有与不带密码学审计的系统相同的可伸缩性,并把性能损失尽可能压低,还需要投入一些工作;但这些技术仍然值得关注。
本章小结
本章以流处理理念为基础,讨论了设计数据系统的新方法。我们从一个观察出发:没有任何一种工具能够高效服务所有可能的用例,因此应用必然需要组合多种不同的软件来实现目标。我们讨论了如何利用批处理与事件流,让数据变更在不同系统之间流动,从而解决这一 数据集成 问题。
在这种方法中,某些系统被指定为权威记录系统,其他数据则通过转换衍生自它们。这样,我们就能维护索引、物化视图、机器学习模型、统计摘要等。衍生与转换过程保持异步和松散耦合,可以防止一处的问题扩散到系统中无关的部分,从而增强整个系统的稳健性与容错能力。
把数据流表示为从一个数据集到另一个数据集的转换,也有助于应用演化:如果要修改某个处理步骤——例如改变索引或缓存的结构——只需让新的转换代码重新处理整个输入数据集,再次衍生出输出。同样,如果出现问题,也可以修复代码并重新处理数据来恢复。
这些过程与数据库内部既有的工作十分相似,因此,我们把数据流应用重新表述为对数据库组件的 分拆,并通过组合这些松散耦合的组件来构建应用。
观察底层数据的变更,就能更新衍生状态;下游消费者还可以继续观察衍生状态本身。我们甚至可以让这种数据流一路抵达显示数据的最终用户设备,从而构建能够动态更新、反映数据变化,并且离线时仍可工作的用户界面。
接下来,我们讨论了如何保证所有这些处理在发生故障时仍然正确。通过异步处理事件、使用端到端请求标识符使操作幂等,并异步检查约束,就能以可伸缩的方式实现强完整性保证。客户端可以等待检查通过,也可以不等待便继续执行,但要承担约束遭到违反、事后必须道歉的风险。这种方法比使用分布式事务的传统方法更可伸缩、更稳健,也更符合许多业务流程的实际运作方式。
围绕数据流构建应用并异步检查约束,可以避免大部分协调,创建既能维护完整性、又有良好性能的系统,即使处于地理分布式环境或发生故障也不例外。最后,我们简要讨论了如何通过审计验证数据完整性、发现损坏,并指出区块链所用的技术也与事件驱动系统颇为相似。
脚注
参考文献
Rachid Belaid. Postgres Full-Text Search is Good Enough! rachbelaid.com, July 2015. Archived at perma.cc/ZVP9-YDCB ↩︎ ↩︎
Philippe Ajoux, Nathan Bronson, Sanjeev Kumar, Wyatt Lloyd, and Kaushik Veeraraghavan. Challenges to Adopting Stronger Consistency at Scale. At 15th USENIX Workshop on Hot Topics in Operating Systems (HotOS), May 2015. ↩︎ ↩︎
Pat Helland and Dave Campbell. Building on Quicksand. At 4th Biennial Conference on Innovative Data Systems Research (CIDR), January 2009. ↩︎ ↩︎ ↩︎ ↩︎ ↩︎ ↩︎
Jessica Kerr. Provenance and Causality in Distributed Systems. jessitron.com, September 2016. Archived at perma.cc/DTD2-F8ZM ↩︎ ↩︎ ↩︎
Jay Kreps. The Log: What Every Software Engineer Should Know About Real-Time Data’s Unifying Abstraction. engineering.linkedin.com, December 2013. Archived at perma.cc/2JHR-FR64 ↩︎ ↩︎ ↩︎ ↩︎
Pat Helland. Life Beyond Distributed Transactions: An Apostate’s Opinion. At 3rd Biennial Conference on Innovative Data Systems Research (CIDR), January 2007. ↩︎ ↩︎
Lionel A. Smith. The Broad Gauge Story. Journal of the Monmouthshire Railway Society, Summer 1985. Archived at perma.cc/DDK9-JA6X ↩︎
Jacqueline Xu. Online Migrations at Scale. stripe.com, February 2017. Archived at perma.cc/ZQY2-EAU2 ↩︎
Flavio Santos and Robert Stephenson. Changing the Wheels on a Moving Bus — Spotify’s Event Delivery Migration. engineering.atspotify.com, October 2021. Archived at perma.cc/5C4V-G8EV ↩︎
Molly Bartlett Dishman and Martin Fowler. Agile Architecture. At O’Reilly Software Architecture Conference, March 2015. ↩︎
Nathan Marz and James Warren. Big Data: Principles and Best Practices of Scalable Real-Time Data Systems. Manning, 2015. ISBN: 978-1-617-29034-3 ↩︎
Jay Kreps. Questioning the Lambda Architecture. oreilly.com, July 2014. Archived at perma.cc/PGH6-XUCH ↩︎ ↩︎
Raul Castro Fernandez, Peter Pietzuch, Jay Kreps, Neha Narkhede, Jun Rao, Joel Koshy, Dong Lin, Chris Riccomini, and Guozhang Wang. Liquid: Unifying Nearline and Offline Big Data Integration. At 7th Biennial Conference on Innovative Data Systems Research (CIDR), January 2015. ↩︎
Dennis M. Ritchie and Ken Thompson. The UNIX Time-Sharing System. Communications of the ACM, volume 17, issue 7, pages 365–375, July 1974. doi:10.1145/361011.361061 ↩︎ ↩︎
Wes McKinney. The Road to Composable Data Systems: Thoughts on the Last 15 Years and the Future. wesmckinney.com, September 2023. Archived at perma.cc/J9SJ-886N ↩︎
Eric A. Brewer and Joseph M. Hellerstein. CS262a: Advanced Topics in Computer Systems. Lecture notes, University of California, Berkeley, cs.berkeley.edu, August 2011. Archived at perma.cc/TE79-LGWU ↩︎
Michael Stonebraker. The Case for Polystores. wp.sigmod.org, July 2015. Archived at perma.cc/G7J2-KR45 ↩︎ ↩︎
Jennie Duggan, Aaron J. Elmore, Michael Stonebraker, Magda Balazinska, Bill Howe, Jeremy Kepner, Sam Madden, David Maier, Tim Mattson, and Stan Zdonik. The BigDAWG Polystore System. ACM SIGMOD Record, volume 44, issue 2, pages 11–16, June 2015. doi:10.1145/2814710.2814713 ↩︎
David B. Lomet, Alan Fekete, Gerhard Weikum, and Mike Zwilling. Unbundling Transaction Services in the Cloud. At 4th Biennial Conference on Innovative Data Systems Research (CIDR), January 2009. ↩︎
Martin Kleppmann and Jay Kreps. Kafka, Samza and the Unix Philosophy of Distributed Data. IEEE Data Engineering Bulletin, volume 38, issue 4, pages 4–14, December 2015. ↩︎
John Hugg. Winning Now and in the Future: Where Volt Active Data Shines. voltactivedata.com, March 2016. Archived at perma.cc/44MP-3MWM ↩︎
Felienne Hermans. Spreadsheets Are Code. At Code Mesh, November 2015. ↩︎
Dan Bricklin and Bob Frankston. VisiCalc: Information from Its Creators. danbricklin.com. Archived at archive.org ↩︎
D. Sculley, Gary Holt, Daniel Golovin, Eugene Davydov, Todd Phillips, Dietmar Ebner, Vinay Chaudhary, and Michael Young. Machine Learning: The High-Interest Credit Card of Technical Debt. At NIPS Workshop on Software Engineering for Machine Learning (SE4ML), December 2014. Archived at https://perma.cc/M3MD-U7WL ↩︎
Peter Bailis, Alan Fekete, Michael J. Franklin, Ali Ghodsi, Joseph M. Hellerstein, and Ion Stoica. Feral Concurrency Control: An Empirical Investigation of Modern Application Integrity. At ACM International Conference on Management of Data (SIGMOD), June 2015. doi:10.1145/2723372.2737784 ↩︎ ↩︎
Guy Steele. Re: Need for Macros (Was Re: Icon). email to ll1-discuss mailing list, people.csail.mit.edu, December 2001. Archived at perma.cc/K9X8-CJ65 ↩︎
Ben Stopford. Microservices in a Streaming World. At QCon London, March 2016. ↩︎ ↩︎
Adam Bellemare. Building Event-Driven Microservices, 2nd Edition. O’Reilly Media, 2025. ↩︎
Christian Posta. Why Microservices Should Be Event Driven: Autonomy vs Authority. blog.christianposta.com, May 2016. Archived at perma.cc/E6N9-3X92 ↩︎
Alex Feyerke. Designing Offline-First Web Apps. alistapart.com, December 2013. Archived at perma.cc/WH7R-S2DS ↩︎
Martin Kleppmann. Turning the Database Inside-out with Apache Samza. at Strange Loop, September 2014. Archived at perma.cc/U6E8-A9MT ↩︎ ↩︎
Sebastian Burckhardt, Daan Leijen, Jonathan Protzenko, and Manuel Fähndrich. Global Sequence Protocol: A Robust Abstraction for Replicated Shared State. At 29th European Conference on Object-Oriented Programming (ECOOP), July 2015. doi:10.4230/LIPIcs.ECOOP.2015.568 ↩︎ ↩︎
Evan Czaplicki and Stephen Chong. Asynchronous Functional Reactive Programming for GUIs. At 34th ACM SIGPLAN Conference on Programming Language Design and Implementation (PLDI), June 2013. doi:10.1145/2491956.2462161 ↩︎
Eno Thereska, Damian Guy, Michael Noll, and Neha Narkhede. Unifying Stream Processing and Interactive Queries in Apache Kafka. confluent.io, October 2016. Archived at perma.cc/W8JG-EAZF ↩︎
Frank McSherry. Dataflow as Database. github.com, July 2016. Archived at perma.cc/384D-DUFH ↩︎
Peter Alvaro. I See What You Mean. At Strange Loop, September 2015. ↩︎
Nathan Marz. Trident: A High-Level Abstraction for Realtime Computation. blog.x.com, August 2012. Archived at archive.org ↩︎
Edi Bice. Low Latency Web Scale Fraud Prevention with Apache Samza, Kafka and Friends. At Merchant Risk Council MRC Vegas Conference, March 2016. Archived at perma.cc/T3H5-QN3R ↩︎
Charity Majors. The Accidental DBA. charity.wtf, October 2016. Archived at perma.cc/6ANP-ARB6 ↩︎
Arthur J. Bernstein, Philip M. Lewis, and Shiyong Lu. Semantic Conditions for Correctness at Different Isolation Levels. At 16th International Conference on Data Engineering (ICDE), February 2000. doi:10.1109/ICDE.2000.839387 ↩︎
Sudhir Jorwekar, Alan Fekete, Krithi Ramamritham, and S. Sudarshan. Automating the Detection of Snapshot Isolation Anomalies. At 33rd International Conference on Very Large Data Bases (VLDB), September 2007. ↩︎
Kyle Kingsbury. Jespen: Distributed Systems Safety Research. jepsen.io. ↩︎
Michael Jouravlev. Redirect After Post. theserverside.com, August 2004. Archived at archive.org ↩︎
Jerome H. Saltzer, David P. Reed, and David D. Clark. End-to-End Arguments in System Design. ACM Transactions on Computer Systems, volume 2, issue 4, pages 277–288, November 1984. doi:10.1145/357401.357402 ↩︎ ↩︎
Peter Bailis, Alan Fekete, Michael J. Franklin, Ali Ghodsi, Joseph M. Hellerstein, and Ion Stoica. Coordination Avoidance in Database Systems. Proceedings of the VLDB Endowment, volume 8, issue 3, pages 185–196, November 2014. doi:10.14778/2735508.2735509 ↩︎ ↩︎
Alex Yarmula. Strong Consistency in Manhattan. blog.x.com, March 2016. Archived at archive.org ↩︎
Martin Kleppmann, Alastair R. Beresford, and Boerge Svingen. Online Event Processing: Achieving consistency where distributed transactions have failed. Communications of the ACM, volume 62, issue 5, pages 43-49, May 2019. doi:10.1145/3312527 ↩︎
Jim Gray. The Transaction Concept: Virtues and Limitations. At 7th International Conference on Very Large Data Bases (VLDB), September 1981. Archived at perma.cc/8VPT-N5H6 ↩︎
Hector Garcia-Molina and Kenneth Salem. Sagas. At ACM International Conference on Management of Data (SIGMOD), May 1987. doi:10.1145/38713.38742 ↩︎
Annamalai Gurusami and Daniel Price. Bug #73170: Duplicates in Unique Secondary Index Because of Fix of Bug#68021. bugs.mysql.com, July 2014. Archived at perma.cc/P6BV-W7JJ ↩︎
Gary Fredericks. Postgres Serializability Bug. github.com, September 2015. Archived at perma.cc/N8UP-2822 ↩︎
Xiao Chen. HDFS DataNode Scanners and Disk Checker Explained. blog.cloudera.com, December 2016. Archived at perma.cc/6S36-X98L ↩︎
Daniel Persson. How does Ceph scrubbing work? youtube.com, March 2022. ↩︎
Jay Kreps. Getting Real About Distributed System Reliability. blog.empathybox.com, March 2012. Archived at perma.cc/9B5Q-AEBW ↩︎
Martin Fowler. The LMAX Architecture. martinfowler.com, July 2011. Archived at perma.cc/5AV4-N6RJ ↩︎
Sam Stokes. Move Fast with Confidence. five-eights.com, July 2016. Archived at perma.cc/J8C6-DHXB ↩︎
Ralph C. Merkle. A Digital Signature Based on a Conventional Encryption Function. At CRYPTO ‘87, August 1987. doi:10.1007/3-540-48184-2_32 ↩︎
Ben Laurie. Certificate Transparency. ACM Queue, volume 12, issue 8, pages 10-19, August 2014. doi:10.1145/2668152.2668154 ↩︎
Mark D. Ryan. Enhanced Certificate Transparency and End-to-End Encrypted Mail. At Network and Distributed System Security Symposium (NDSS), February 2014. doi:10.14722/ndss.2014.23379 ↩︎