12 流处理

有效的复杂系统总是从简单的系统演化而来。反之亦然:从零设计的复杂系统没一个能有效工作的。
—— 约翰・加尔,Systemantics(1975)
在 第 11 章 中,我们讨论了批处理技术:它以一组文件作为输入,并生成一组新的输出文件。输出是 衍生数据(derived data) 的一种形式;也就是说,如有必要,可以再次运行批处理来重新创建这份数据集。我们看到,这个简单而强大的思想可以用来构建搜索索引、推荐系统和分析系统,等等。
然而,第 11 章 始终建立在一个重要假设之上:输入是 有界的(bounded),也就是大小已知且有限,因此批处理知道何时已经读完输入。例如,作为 MapReduce 核心环节的排序操作必须先读完全部输入,才能开始生成输出。因为最后一条输入记录可能恰好拥有最小的键,因而需要成为第一条输出记录,所以不能提前开始输出。
实际上,许多数据都是 无界的(unbounded),因为它们会随时间推移陆续到达:用户昨天和今天产生了数据,明天还会继续产生更多数据。只要你的企业还在经营,这个过程就不会结束,因此从任何有意义的角度来看,数据集都永远不会“完整”1。所以,批处理程序不得不人为地按固定时长切分数据,例如每天结束时处理当天的数据,或每小时结束时处理这一小时的数据。
每日批处理的问题在于,输入中的变化要到一天后才会反映到输出中,对许多没有耐心的用户来说实在太慢。为了缩短延迟,可以更频繁地运行处理——比如每秒结束时处理这一秒的数据——也可以完全抛开固定的时间切片,改为连续处理,每个事件一发生就立即处理。这就是 流处理(stream processing) 的基本思想。
一般来说,“流”是指随时间推移逐渐可用的数据。这个概念出现在许多地方:Unix 的 stdin 和 stdout、编程语言中的惰性列表2、文件系统 API(例如 Java 的 FileInputStream)、TCP 连接、通过互联网传输的音频和视频,等等。
本章将把 事件流(event stream) 作为一种数据管理机制来考察:它是上一章批量数据的无界、增量处理版本。我们首先讨论如何表示、存储流,以及如何通过网络传输流;接着在“数据库与流”中研究流与数据库的关系;最后在“流处理”中探讨持续处理这些流的方法和工具,以及如何用它们构建应用。
传递事件流
在批处理领域,作业的输入和输出都是文件(可能位于分布式文件系统上)。那么,流处理领域中的对应物是什么?
当输入是文件(即字节序列)时,第一个处理步骤通常是把它解析成一系列记录。在流处理的语境中,记录更常被称为 事件(event),但两者本质上是同一种东西:一个小巧、自包含、不可变的对象,记录某个时刻发生的某件事。事件通常带有时间戳,表示按照日历时钟来看它发生在何时(参见“单调时钟与日历时钟”)。
例如,事件可能来自用户的某个操作,如浏览页面或完成购买;也可能来自机器,如温度传感器的定期测量值或 CPU 利用率指标。在“使用 Unix 工具的批处理”示例中,Web 服务器日志的每一行就是一个事件。
事件可以编码成文本字符串、JSON,或 第 5 章 讨论的某种二进制格式。编码后,你就可以存储事件,例如把它追加到文件、插入关系表,或写入文档数据库;也可以通过网络把事件发送到另一个节点进行处理。
在批处理中,文件写入一次,之后可能由多个作业读取。与之类似,在流处理术语中,事件由 生产者(producer)(也称为 发布者(publisher) 或 发送者(sender))生成一次,之后可能由多个 消费者(consumer)(也称为 订阅者(subscriber) 或 接收者(recipient))处理3。在文件系统中,文件名标识一组相关记录;在流式系统中,相关事件通常被归入同一个 主题(topic) 或 流(stream)。
原则上,文件或数据库足以把生产者和消费者连接起来:生产者将生成的每个事件写入数据存储,而各个消费者定期轮询数据存储,检查自上次运行以来出现了哪些新事件。这实质上就是每日结束时处理当天数据的批处理程序所做的事情。
然而,当我们转向低延迟的持续处理时,如果数据存储并非为这种用法而设计,轮询的代价就会很高。轮询越频繁,返回新事件的请求比例越低,额外开销也就越大。更好的做法是在出现新事件时通知消费者。
传统数据库对这种通知机制的支持并不好。关系数据库通常提供 触发器(trigger),可以对变更作出响应(例如向表中插入一行),但触发器能做的事情非常有限,在数据库设计中多少有些事后补充的意味4。因此,人们开发了专门用于传递事件通知的工具。
消息传递系统
向消费者通知新事件的一种常见方法是使用 消息传递系统(messaging system):生产者发送包含事件的消息,再由系统将消息推送给消费者。我们之前在“事件驱动的架构”中提到过这类系统,现在来进一步了解其细节。
在生产者与消费者之间建立 Unix 管道或 TCP 连接等直接通信信道,是实现消息传递系统的一种简单方法。不过,大多数消息传递系统都扩展了这个基本模型。Unix 管道和 TCP 恰好连接一个发送者与一个接收者,而消息传递系统允许多个生产者节点向同一主题发送消息,也允许多个消费者节点接收一个主题中的消息。
在这种 发布/订阅 模型中,不同系统采取了五花八门的方法,没有一种答案适合所有用途。要区分这些系统,下面两个问题尤其有用:
如果生产者发送消息的速度超过消费者的处理速度,会发生什么? 大体上有三种选择:系统可以丢弃消息、把消息缓存在队列中,或施加 背压(backpressure)(也称为 流量控制(flow control),即阻塞生产者,令其暂停发送更多消息)。例如,Unix 管道和 TCP 都采用背压:它们有一个固定大小的小缓冲区;如果缓冲区已满,发送者就会被阻塞,直到接收者从中取走数据(参见“网络拥塞与排队”)。
如果消息被缓存在队列中,就必须弄清楚队列不断增长会有什么后果:队列大到内存容纳不下时,系统会崩溃,还是会把消息写入磁盘?如果写入磁盘,磁盘访问会怎样影响消息传递系统的性能5?磁盘写满时又会发生什么6?
如果节点崩溃或暂时离线,会发生什么——是否会丢失消息? 和数据库一样,要实现持久性,可能需要以某种方式组合使用磁盘写入与复制(参见侧栏“复制与持久性”),而这会付出代价。如果能够容忍偶尔丢失消息,那么在同样的硬件上,通常可以获得更高的吞吐量和更低的延迟。
能否容忍消息丢失,很大程度上取决于应用。例如,对于定期发送的传感器读数和指标,偶尔缺失一个数据点或许并不重要,因为不久后还会发来更新值。不过要小心:如果大量消息被丢弃,你可能无法立即察觉指标已经失真7。如果你要对事件计数,可靠传递就重要得多,因为每丢失一条消息,计数器就会出现一份误差。
我们在 第 11 章 探讨的批处理系统有一个很好的特性:它们提供了很强的可靠性保证。失败的任务会自动重试,失败任务产生的部分输出会自动丢弃。因此,最终输出与从未发生故障时相同,这简化了编程模型。本章后面将考察如何在流处理环境中提供类似的保证。
直接从生产者传递给消费者
许多消息传递系统让生产者与消费者直接通过网络通信,不经过任何中间节点:
UDP 组播广泛用于金融行业中的股票行情等数据流,因为这些场景十分看重低延迟8。尽管 UDP 本身并不可靠,但应用层协议可以恢复丢失的数据包(生产者必须记住已发送的数据包,以便按需重传)。
ZeroMQ 和 nanomsg 等无代理消息库采用类似的方法,通过 TCP 或 IP 组播实现发布/订阅消息传递。
StatsD 等指标收集代理9使用不可靠的 UDP 消息,从网络中的所有机器收集指标并进行监控。(在 StatsD 协议中,只有收到全部消息,计数器指标才是准确的;使用 UDP 意味着这些指标至多只能是近似值10。另见“TCP 与 UDP”。)
如果消费者在网络上公开了一项服务,生产者可以直接发起 HTTP 或 RPC 请求(参见“流经服务的数据流:REST 与 RPC”),把消息推送给消费者。Webhook11就建立在这一思想之上:把一个服务的回调 URL 注册到另一个服务中,后者每逢事件发生便向该 URL 发出请求。
尽管这些直接消息传递系统在各自的目标场景中表现良好,但应用代码通常必须自行考虑消息丢失的可能性。它们能容忍的故障十分有限:即使协议可以检测并重传网络中丢失的数据包,通常仍会假设生产者和消费者始终在线。
消费者离线时,可能会错过其不可达期间发送的消息。有些协议允许生产者重试失败的消息传递,但如果生产者崩溃,丢失了本应重试的消息缓冲区,这种办法就可能失效。
消息代理
一种广泛使用的替代方案是通过 消息代理(message broker)(也称为 消息队列(message queue))发送消息。消息代理实质上是一类针对消息流优化的数据库12。它作为服务器运行,生产者和消费者则作为客户端连接到它。生产者把消息写入代理,消费者通过读取代理来接收消息。
数据集中到代理之后,系统更容易容忍客户端时而连接、时而断开乃至崩溃,持久性问题也转移给了代理。有些消息代理只把消息保存在内存中,另一些则会根据配置把消息写入磁盘,以免代理崩溃时丢失消息。面对缓慢的消费者,它们通常允许队列无限增长,而不是丢弃消息或施加背压,不过具体行为也可能取决于配置。
排队还有一个后果:消费者通常是 异步的。生产者发送消息时,一般只等待代理确认消息已被缓存,并不等待消费者完成处理。消息会在未来某个无法确定的时间点传递给消费者——往往不到一秒,但如果队列中积压了大量消息,也可能晚得多。
消息代理与数据库的对比
有些消息代理甚至可以通过 XA 或 JTA 参与两阶段提交协议(参见“跨不同系统的分布式事务”)。这一功能使它们在性质上与数据库十分相似,不过消息代理与数据库之间仍有一些重要的实际差异:
数据库通常会一直保存数据,直到有人显式删除;而有些消息代理会在消息成功传递给消费者后自动将其删除。这类消息代理不适合长期存储数据。
正因为消息会很快删除,大多数消息代理假定工作集相当小,也就是队列很短。如果消费者缓慢,导致代理必须缓冲大量消息(内存容纳不下时还可能溢写到磁盘),每条消息的处理时间就会延长,总体吞吐量也可能下降5。
数据库通常支持二级索引,并能通过查询语言以多种方式搜索数据;消息代理通常只支持订阅与某种模式匹配的一组主题。两者本质上都让客户端选择自己关心的数据子集,但数据库提供的查询功能通常强大得多。
查询数据库时,结果通常以某一时刻的数据快照为依据;如果另一个客户端随后写入数据库并改变了查询结果,第一个客户端不会知道先前的结果已经过时,除非再次执行查询或轮询变更。相比之下,消息代理不支持任意查询,也不允许修改已经发出的消息,但会在数据发生变化时(即有新消息可用时)通知客户端。
这是消息代理的传统形态,JMS13、AMQP14 等标准对它作出了规范,RabbitMQ、ActiveMQ、HornetQ、Qpid、TIBCO Enterprise Message Service、IBM MQ、Azure Service Bus 和 Google Cloud Pub/Sub 等软件则实现了这种形态15。尽管也可以把数据库用作队列,但要把性能调到理想水平并不容易16。
多个消费者
多个消费者读取同一主题中的消息时,主要有两种消息传递模式,如 图 12-1 所示:
- 负载均衡
每条消息只传递给 一个 消费者,因此多个消费者可以分担处理该主题消息的工作。代理可以任意把消息分配给消费者。如果消息处理成本很高,而你希望通过增加消费者来并行处理,便适合采用这种模式。(在 AMQP 中,可以让多个客户端消费同一个队列来实现负载均衡;在 JMS 中,这称为 共享订阅(shared subscription)。)
- 扇出
每条消息都传递给 所有 消费者。扇出让几个相互独立的消费者都能“收听”同一份消息广播,彼此互不影响——相当于流处理版本的“几个不同批处理作业读取同一个输入文件”。(JMS 的主题订阅和 AMQP 的交换器绑定提供了这一功能。)

这两种模式可以结合使用,例如 Kafka 的 消费者组(consumer group) 功能。消费者组订阅某个主题后,该主题中的每条消息都会发送给组内的一个消费者(在组内消费者之间进行负载均衡)。如果两个不同的消费者组订阅同一主题,那么每条消息都会发送给各组中的一个消费者(在消费者组之间实现扇出)。
确认应答与重新传递
消费者随时可能崩溃,因此代理可能已经把消息传递给消费者,但消费者尚未处理,或只处理了一部分便发生崩溃。为了确保消息不会丢失,消息代理采用 确认应答(acknowledgment):客户端处理完消息后,必须明确告知代理,代理才能将消息从队列中移除。
如果客户端连接关闭或超时,而代理尚未收到确认应答,代理便假定消息没有得到处理,并把它重新传递给另一个消费者。(请注意,消息可能 实际上已经 处理完毕,只是确认应答在网络中丢失了。除非该操作具有幂等性,或不要求恰好一次语义,否则必须使用原子提交协议来处理这种情况,详见“恰好一次消息处理”。)
重新传递与负载均衡结合后,会对消息顺序产生一个有趣的影响。在 图 12-2 中,消费者通常按照生产者发送消息的顺序进行处理。但是,消费者 2 在处理消息 m3 时崩溃,与此同时消费者 1 正在处理消息 m4。尚未确认的消息 m3 随后被重新传递给消费者 1,于是消费者 1 按照 m4、m3、m5 的顺序处理消息。因此,m3 和 m4 的传递顺序与生产者 1 的发送顺序不同。

即使消息代理试图保持消息顺序(JMS 和 AMQP 标准都作此要求),负载均衡与重新传递的结合仍不可避免地会使消息发生重排。要避免这个问题,可以为每个消费者使用单独的队列,也就是不使用负载均衡。如果消息彼此完全独立,重排并无大碍;但如果消息之间存在因果依赖,顺序就可能十分重要,本章后面将看到这一点。
重新传递还可能浪费资源、造成资源饥饿,甚至永久阻塞一条流。常见情形是生产者没有正确序列化消息,例如编码成 JSON 的对象缺少必填键。任何读到这条消息的消费者都会期待该键存在,并在发现缺失时失败。由于没有发出确认应答,代理会再次发送消息,使另一个消费者也随之失败,如此无限循环。如果代理提供严格的顺序保证,后续处理将完全无法推进。允许消息重排的代理仍能继续处理其他消息,却会把资源浪费在永远得不到确认的消息上。
死信队列(dead letter queue,DLQ) 用来处理这类问题。系统不再把消息留在当前队列中无限重试,而是将其移到另一条队列,让消费者能够继续前进17、18。通常会对死信队列设置监控——队列中出现任何消息都意味着发生了错误。检测到新消息后,操作员可以决定永久丢弃它、手动修改后重新生产这条消息,或修复消费者代码,使其能够正确处理消息。大多数队列系统都提供 DLQ;如今,Apache Pulsar 等基于日志的消息系统以及 Kafka Streams 等流处理系统也开始支持它19。
基于日志的消息代理
通过网络发送数据包或向网络服务发出请求,通常都是转瞬即逝的操作,不会留下永久痕迹。尽管可以通过抓包和日志记录把这些操作永久保存下来,但我们通常不会这样看待它们。AMQP/JMS 风格的消息代理继承了这种临时消息传递的思路:即使把消息写入磁盘,也会在消息传递给消费者后很快将其删除。
数据库和文件系统的思路恰恰相反:凡是写入数据库或文件的内容,通常都应该永久保存,至少要保存到有人明确决定再次删除它为止。
这种思路上的差异会极大地影响衍生数据的创建方式。正如 第 11 章 所说,批处理的一项关键特性是可以反复运行,试验不同的处理步骤,又不用担心损坏输入,因为输入是只读的。AMQP/JMS 风格的消息传递却并非如此:如果确认应答会让代理删除消息,那么接收消息就是一种破坏性操作。你无法重新运行同一个消费者,并期望得到同样的结果。
向消息传递系统加入新消费者时,它通常只能接收注册之后发出的消息;先前的消息早已消失,无法恢复。相比之下,文件和数据库可以随时加入新客户端,并读取任意久远的数据,只要应用没有明确覆盖或删除这些数据。
为什么不能把两者结合起来,既采用数据库的持久存储方式,又具备消息传递的低延迟通知能力?这就是 基于日志的消息代理(log-based message broker) 的思想。近年来,这类系统已经十分流行。
使用日志进行消息存储
日志就是磁盘上一系列仅追加的记录。我们此前在 第 4 章 讨论日志结构存储引擎和预写日志时、在 第 6 章 讨论复制时,以及在 第 10 章 把日志作为一种共识形式来讨论时,都见过这种结构。
同一种结构也可以用来实现消息代理:生产者把消息追加到日志末尾,以此发送消息;消费者顺序读取日志,以此接收消息。消费者读到日志末尾后,便等待新消息追加的通知。用于监视文件新增内容的 Unix 工具 tail -f,实质上的工作方式与此相同。
为了把吞吐量扩展到单块磁盘的能力之上,可以对日志进行 分片(即 第 7 章 所说的分片)。不同分片可以托管在不同机器上,每个分片都是一份独立于其他分片读写的日志。一个主题可以定义为一组承载相同类型消息的分片。图 12-3 展示了这种方法。
在每个分片中——Kafka 把它称为 分区(partition)——代理会为每条消息分配一个单调递增的序列号,也就是 偏移量(offset)(图 12-3 方框中的数字就是消息偏移量)。分区是仅追加的,因此这样的序列号有明确含义:分区内的消息具有全序,而不同分区之间没有顺序保证。

Apache Kafka20 和 Amazon Kinesis Streams 都是以这种方式工作的基于日志的消息代理。Google Cloud Pub/Sub 的架构与之相似,但公开的是 JMS 风格的 API,而不是日志抽象15。尽管这些消息代理会把所有消息写入磁盘,但通过跨多台机器分片,仍能达到每秒数百万条消息的吞吐量,并通过复制消息实现容错21、22。
日志与传统消息传递的比较
基于日志的方法自然支持扇出式消息传递,因为多个消费者可以各自读取日志而互不影响——读取消息不会把它从日志中删除。要在一组消费者之间实现负载均衡,代理可以把整个分片分配给消费者组中的节点,而不是把一条条消息分别分配给各个消费者客户端。
随后,每个客户端会消费所分配分片中的 全部 消息。消费者得到一个日志分片后,通常会采用直截了当的单线程方式,顺序读取其中的消息。这种粗粒度负载均衡有一些缺点:
分担主题消费工作的节点数,最多只能等于该主题的日志分片数,因为同一分片内的消息都传递给同一个节点。(也可以设计一种让两个消费者共同处理一个分片的负载均衡方案:两者都读取完整消息集,但一个只处理偏移量为偶数的消息,另一个只处理偏移量为奇数的消息。另一种办法是把消息处理分派给线程池,但这会使消费者偏移量的管理更加复杂。一般来说,最好还是用单线程处理一个分片,并通过增加分片来提高并行度。)
如果某一条消息处理得很慢,就会阻塞该分片中后续消息的处理,这是一种队头阻塞(参见“描述性能”)。
因此,如果消息处理成本很高,你希望以消息为单位并行处理,而且消息顺序并不十分重要,那么 JMS/AMQP 风格的消息代理更合适。反之,如果消息吞吐量很高,每条消息都能迅速处理,而且顺序非常重要,基于日志的方法就表现出色23、24。不过,两种架构之间的界线正在变得模糊:Kafka 等基于日志的消息系统如今也支持 JMS/AMQP 风格的消费者组,让多个消费者能够接收同一分区中的消息25、26。
由于分片日志通常只能保持单个分片内部的消息顺序,所有必须以一致顺序处理的消息都需要路由到同一个分片。例如,应用可能要求与某个特定用户有关的事件始终以固定顺序出现。可以根据事件的用户 ID 选择分片来实现这一点;换句话说,把用户 ID 作为 分区键(partition key)。
消费者偏移量
顺序消费一个分片,很容易判断哪些消息已经处理:偏移量小于消费者当前偏移量的消息都已处理,偏移量更大的消息则尚未读到。因此,代理无需跟踪每条消息的确认应答,只需定期记录消费者偏移量。这样既减少了簿记开销,也带来了批量处理和流水线化的机会,有助于提高基于日志系统的吞吐量。不过,如果消费者发生故障,它会从上次记录的偏移量恢复,而不是从自己实际读到的最新位置恢复,因此可能会再次读到一些消息。
这种偏移量其实与单主数据库复制中常见的 日志序列号 非常相似,我们曾在“设置新的副本”中讨论过它。在数据库复制中,日志序列号允许追随者断开连接后重新连接到领导者,并且不跳过任何写入就恢复复制。这里采用的原理完全相同:消息代理扮演领导者数据库的角色,消费者则像追随者。
如果消费者节点发生故障,消费者组会把它的分片分配给另一个节点,后者从最后记录的偏移量开始消费。如果原消费者已经处理了后续消息,却尚未记录相应的偏移量,那么重启后这些消息会被再次处理。本章后面将讨论如何应对这一问题。
磁盘空间使用
如果永远只向日志追加内容,磁盘空间终有耗尽之时。为了回收空间,日志实际上会被切成若干段,旧段会不时被删除或移入归档存储。(我们将在“日志压实”中讨论一种更复杂的空间回收方法。)
这意味着,如果缓慢的消费者跟不上消息产生的速度,落后到其偏移量指向已经删除的日志段,就会错过一些消息。实际上,日志实现了一个大小有限的缓冲区,填满后便丢弃旧消息,也就是 循环缓冲区(circular buffer) 或 环形缓冲区(ring buffer)。不过,由于这个缓冲区位于磁盘上,它可以相当大。
让我们粗略估算一下。在本书写作时,典型的大容量硬盘为 20 TB,顺序写入吞吐量为 250 MB/s。如果始终以最高速度写入消息,大约 22 小时后磁盘就会写满,必须开始删除最旧的消息。这意味着,即使有许多机器和许多磁盘,磁盘日志也总能缓冲至少 22 小时的消息,因为增加磁盘会同时增加可用空间和总写入带宽。实际部署很少用满磁盘的全部写入带宽,所以日志通常可以保存数天乃至数周的消息。
许多基于日志的消息代理如今会把消息存入对象存储,以扩大存储容量,这与“以对象存储为后端的数据库”中讨论的做法相似。Apache Kafka 和 Redpanda 等消息代理通过分层存储从对象存储提供较旧的消息;WarpStream、Confluent Freight 和 Bufstream 等系统则把全部数据都存入对象存储。除了成本效益,这种架构也简化了数据集成:对象存储中的消息以 Iceberg 表保存,批处理作业和数据仓库作业可以直接在这些数据上运行,无需先把数据复制到另一个系统。
当消费者跟不上生产者时
在“消息传递系统”开头,我们讨论了消费者跟不上生产者发送消息的速度时可采取的三种办法:丢弃消息、缓冲消息或施加背压。按照这种分类,基于日志的方法属于缓冲,它提供了一个很大但大小固定的缓冲区,其上限由可用磁盘空间决定。
如果消费者远远落后,需要的消息已经早于磁盘保留范围,它就无法再读取这些消息——也就是说,代理实际上会丢弃超出缓冲容量的旧消息。你可以监控消费者落后日志头部多远,并在落后过多时发出告警。由于缓冲区很大,通常有足够时间让运维人员修复缓慢的消费者,并在它开始漏掉消息前追上进度。
即使某个消费者确实落后太多并开始漏掉消息,也只有它自己受到影响,不会干扰其他消费者的服务。这是一项很大的运维优势:你可以出于开发、测试或调试目的,试验性地消费生产日志,不必太担心干扰生产服务。消费者关闭或崩溃后便不再消耗资源,唯一留下的只是它的偏移量。
这种行为也与传统消息代理形成鲜明对比。在传统代理中,必须小心删除消费者已经关闭的队列,否则这些队列会继续积累不再需要的消息,挤占仍在活动的消费者可用的内存。
重播旧消息
前面提到,对于 AMQP/JMS 风格的消息代理,处理并确认消息是一种破坏性操作,因为这会使代理删除消息。而在基于日志的消息代理中,消费消息更像读取文件:它是一种不会改变日志的只读操作。
除了消费者本身产生的输出,处理消息唯一的副作用就是消费者偏移量向前移动。但偏移量由消费者控制,所以必要时很容易调整。例如,可以用昨天的偏移量启动一份消费者副本,把输出写到另一个位置,从而重新处理过去一天的消息。你可以任意重复这个过程,每次换用不同的处理代码。
这一特性让基于日志的消息传递更像上一章的批处理:通过可重复执行的转换过程,把衍生数据与输入数据明确分开。它为试验提供了更多空间,也更容易从错误和程序缺陷中恢复,因此很适合用来集成组织内部的数据流27。
数据库与流
前面我们对消息代理与数据库做过一些比较。传统上,它们被视为两类不同的工具,但基于日志的消息代理已经成功地把数据库中的思想用于消息传递。反过来也一样:我们可以把消息传递和流中的思想用于数据库。
一种做法是用 事件流充当存储数据的权威记录系统(参阅 “权威记录系统与衍生数据”)。这正是我们在 “事件溯源与 CQRS” 中讨论过的 事件溯源(event sourcing):不再用更新和删除来改变数据模型,而是把每次状态变化建模成不可变事件,写入仅追加日志;所有为读取优化的物化视图都从这些事件中衍生出来。基于日志的消息代理采用仅追加存储,还能以低延迟通知消费者有新事件到来,因此很适合事件溯源——只要将其配置为永不删除旧事件。
不过,你不必走到采用事件溯源这一步;即便数据模型是可变的,事件流对数据库仍然很有用。事实上,每次数据库写入都是一个可以捕获、存储和处理的事件。数据库与流的联系不只是日志在磁盘上的物理存储形式,而是更为根本。
例如,复制日志(参阅 “复制日志的实现”)就是数据库写入事件组成的流,由领导者在处理事务时生成。追随者把这股写入流应用到自己的数据库副本上,最终得到同一份数据的准确副本。复制日志中的事件描述的正是已经发生的数据变化。
我们还在 “使用共享日志” 中遇到过 状态机复制(state machine replication)原理:如果每个事件都代表一次数据库写入,并且每个副本都以相同顺序处理相同事件,那么所有副本最终都会达到相同状态(这里假定事件处理是确定性的操作)。这又是事件流的一个例子!
本节先考察异构数据系统中会出现的一个问题,再探讨如何把事件流的思想引入数据库来解决它。
保持系统同步
正如全书反复说明的,没有一个系统能够满足所有数据存储、查询和处理需求。实践中,大多数稍具规模的应用都要组合多种技术才能满足要求:例如用 OLTP 数据库处理用户请求,用缓存加速常见请求,用全文索引处理搜索查询,再用数据仓库进行分析。每个系统都有自己的数据副本,并采用为自身用途优化的表示形式。
相同或相关的数据分散在不同位置,就必须彼此保持同步:数据库中的某个项目更新后,缓存、搜索索引和数据仓库也要随之更新。数据仓库通常通过 ETL 流程完成同步(参阅 “数据仓库”):先取得数据库的完整副本,转换数据,再批量加载进数据仓库——换言之,这是一个批处理过程。同样,我们在 “批处理用例” 中看到,搜索索引、推荐系统以及其他衍生数据系统也可以通过批处理来创建。
如果定期转储完整数据库太慢,有时会改用 双写(dual write):数据发生变化时,应用代码显式写入每个系统,例如先写数据库,再更新搜索索引,最后使相应的缓存项失效(也可能并发执行这些写入)。
然而,双写存在一些严重问题,其中之一就是 图 12-4 所示的竞态条件。在这个例子中,两个客户端并发更新项目 X:客户端 1 想把值设为 A,客户端 2 想把值设为 B。两个客户端都先把新值写入数据库,再写入搜索索引。由于时序不巧,请求交错执行:数据库先收到客户端 1 将值设为 A 的写入,再收到客户端 2 将值设为 B 的写入,因此数据库中的最终值为 B;搜索索引却先收到客户端 2 的写入,再收到客户端 1 的写入,因此最终值为 A。虽然没有发生任何错误,两个系统却永久地不一致了。

除非另外采用并发检测机制,例如我们在 “检测并发写入” 中讨论的版本向量,否则你甚至不会察觉发生过并发写入——一个值只会悄无声息地覆盖另一个值。
双写的另一个问题是,其中一次写入可能失败,另一次却成功。这属于容错问题而不是并发问题,但同样会使两个系统彼此不一致。要保证两次写入要么都成功,要么都失败,就要解决代价高昂的原子提交问题(参阅 “两阶段提交(2PC)”)。
如果只有一个采用单主复制的数据库,那么领导者会决定写入顺序,状态机复制就能在数据库的各个副本之间正常工作。然而,图 12-4 中并不存在唯一的领导者:数据库可能有自己的领导者,搜索索引也可能有自己的领导者,但两者谁也不追随谁,因而可能发生冲突(参阅 “多主复制”)。
如果真能只有一个领导者——例如数据库——并让搜索索引成为数据库的追随者,情况就会好得多。但实践中能做到吗?
变更数据捕获
大多数数据库的复制日志长期以来都被视为内部实现细节,而不是公共 API。客户端理应通过数据库的数据模型和查询语言进行查询,而不是解析复制日志,尝试从中提取数据。
几十年来,许多数据库根本没有提供文档化的方式来获取写入其中的变更日志。因此,要取得数据库中发生的所有变化,再将其复制到搜索索引、缓存或数据仓库等其他存储技术中,一直非常困难。
近年来,变更数据捕获(change data capture,CDC)越来越受关注。它是这样一个过程:观察写入数据库的所有数据变化,将其提取成可以复制到其他系统的形式28。如果变化一经写入就立即以流的形式提供出来,CDC 尤其有用。
例如,你可以捕获数据库中的变化,并持续把相同变化应用到搜索索引。只要按同一顺序应用变更日志,搜索索引中的数据就有望与数据库保持一致。搜索索引以及其他衍生数据系统,都只是变更流的消费者。
图 12-5 展示了 CDC 如何解决 图 12-4 中的并发问题。将 X 分别设为 A 和 B 的两个请求虽然并发到达数据库,但数据库会决定某种执行顺序,并按该顺序把它们写入复制日志。搜索索引再按相同顺序取得并应用这些变化。如果还需要把数据送往数据仓库等其他系统,只须再为 CDC 事件流添加一个消费者。

变更数据捕获的实现
按照 “权威记录系统与衍生数据” 中的说法,我们可以把日志消费者称为 衍生数据系统(derived data system):搜索索引和数据仓库中存储的数据,不过是权威记录系统中数据的另一种视图。变更数据捕获机制确保权威记录系统中的所有变化也会反映到衍生数据系统中,使衍生系统持有准确的数据副本。
实质上,变更数据捕获让一个数据库成为领导者(即从中捕获变化的数据库),让其他系统成为追随者。基于日志的消息代理能够保持消息顺序,避免 图 12-2 中的乱序问题,因此很适合把变更事件从源数据库传送到衍生系统。
逻辑复制日志可以用来实现变更数据捕获(参阅 “逻辑(基于行)的日志复制”),不过需要应对模式变更、恰当建模更新等挑战。开源项目 Debezium 正是为解决这些问题而生。它为 MySQL、PostgreSQL、Oracle、SQL Server、Db2、Cassandra 以及其他许多数据库提供了 源连接器(source connector)。这些连接器接入数据库复制日志,以标准事件模式呈现其中的变化;随后便可转换消息并将其写入下游数据库。Kafka Connect 框架也为各种数据库提供了更多 CDC 连接器。Maxwell 通过解析 binlog 为 MySQL 提供类似功能29;GoldenGate 为 Oracle 提供类似功能;pgcapture 则面向 PostgreSQL。
与消息代理一样,变更数据捕获通常也是异步的:权威记录数据库提交变化之前,并不会等待消费者应用该变化。这种设计在运维上的好处是,增加一个缓慢的消费者不会对权威记录系统造成太大影响;缺点则是复制延迟的所有问题同样存在(参阅 “复制延迟的问题”)。
初始快照
如果拥有数据库有史以来的全部变更日志,就可以通过重播日志来重建数据库的完整状态。然而,永久保留所有变化往往会占用太多磁盘空间,重播也会耗时过久,因此日志通常需要截断。
例如,构建新的全文索引需要整个数据库的完整副本——只应用最近的变更日志还不够,因为其中缺少最近没有更新过的项目。因此,如果没有完整的历史日志,就需要从一个一致的快照开始,正如 “设置新的副本” 中所讨论的那样。
数据库快照必须与变更日志中的某个已知位置或偏移量对应,这样才能知道快照处理完毕后应从何处开始应用变化。有些 CDC 工具集成了快照功能,有些则需要手工完成。Debezium 使用 Netflix 的 DBLog 水位线算法提供增量快照30、31。
日志压实
如果只能保留有限的日志历史,那么每增加一个新的衍生数据系统,都要重新执行一遍快照流程。不过,日志压实(log compaction)提供了一个很好的替代方案。
我们在 “日志结构存储” 中讨论日志结构存储引擎时介绍过日志压实(示例见 图 4-3)。原理很简单:存储引擎定期查找日志中键相同的记录,丢弃重复项,只保留每个键的最新更新。这可能会显著缩小日志段,因此在压实过程中也可以合并日志段,如 图 12-6 所示。整个过程在后台运行。

在日志结构存储引擎中,带有特殊空值的更新(称为 墓碑,tombstone)表示某个键已被删除,并使该键在日志压实时被移除。但只要一个键没有被覆盖或删除,它就会永久留在日志中。压实后的日志所需磁盘空间只取决于数据库当前的内容,而与数据库有史以来发生过多少次写入无关。如果同一个键经常被覆盖,旧值最终会被垃圾回收,只留下最新值。
同样的思路也适用于基于日志的消息代理和变更数据捕获。如果 CDC 系统保证每次变化都有主键,而且对某个键的每次更新都会取代该键的旧值,那么只保留这个键最近一次写入就足够了。
这样一来,每当需要重建搜索索引等衍生数据系统时,都可以让新消费者从经过日志压实的主题的偏移量 0 开始,依次扫描日志中的所有消息。日志保证包含数据库中每个键的最新值(也可能包含一些旧值)——换言之,无须再次对 CDC 源数据库创建快照,便可由此取得数据库内容的完整副本。
Apache Kafka 支持日志压实。正如本章后面将会看到的,它使消息代理不仅可以传送临时消息,还能用作持久存储。
变更流的 API 支持
如今,大多数主流数据库都把变更流作为一等接口公开出来,而不再依赖过去那种事后加装、逆向工程得到的 CDC。MySQL、PostgreSQL 等关系数据库通常通过自身副本所使用的同一份复制日志发送变化。大多数云厂商也为自家产品提供 CDC 方案:例如,Datastream 可以流式访问 Google Cloud 的关系数据库和数据仓库。
即便 Cassandra 这类最终一致、基于法定人数的数据库,如今也支持变更数据捕获。正如 “线性一致性与仲裁” 中所述,客户端必须把写入持久化到多数节点,该写入才被视为可见。法定人数写入很难支持 CDC,因为没有唯一的权威数据源可供订阅;数据是否可见,取决于每个读取者选择的一致性级别。Cassandra 绕开了这个问题:它不提供统一的变更流,而是公开每个节点的原始日志段。消费数据的系统必须读取每个节点的原始日志段,再自行决定如何将其合并成一股流,做法很像法定人数读取者32。
Kafka Connect33 把许多数据库系统的变更数据捕获工具与 Kafka 集成起来。变更事件进入 Kafka 之后,既可以用来更新搜索索引等衍生数据系统,也可以送入本章后面讨论的流处理系统。
变更数据捕获与事件溯源
我们来比较一下变更数据捕获与事件溯源。与变更数据捕获相似,事件溯源也把应用状态的所有变化存成变更事件日志。两者最大的区别在于抽象层次不同:
在变更数据捕获中,应用以可变方式使用数据库,可以随意更新和删除记录。变更日志从数据库底层提取(例如解析复制日志),从而确保提取出的写入顺序与实际写入顺序一致,避免 图 12-4 中的竞态条件。
在事件溯源中,应用逻辑显式构建在写入事件日志的不可变事件之上。事件存储只允许追加,通常不鼓励或禁止更新、删除事件。事件旨在反映应用层发生的事情,而不是底层的状态变化。
哪一种更好取决于具体情况。对于原本没有采用事件溯源的应用,改用事件溯源是一项重大变化,也会带来 “事件溯源与 CQRS” 中讨论的各种利弊。相比之下,CDC 可以用很少的改动接入现有数据库——写入数据库的应用甚至可能根本不知道 CDC 正在运行。
变更数据捕获看起来比事件溯源更容易采用,但它也有自己的一系列挑战。
在微服务架构中,一个数据库通常只由一个服务访问。其他服务通过该服务的公共 API 与之交互,一般不会直接访问数据库。这样,数据库就成为该服务的内部实现细节,开发者可以改变数据库模式而不影响公共 API。
然而,CDC 系统复制数据时通常会沿用上游数据库的模式,这会让这些模式变成公共 API,必须像服务的公共 API 一样加以管理。如果开发者删除数据库表中的一列,依赖该字段的下游消费者就会崩溃。这类挑战一直存在于数据流水线中,但过去通常只影响数据仓库 ETL。CDC 往往以数据流实现,其他生产服务也可能是消费者,因此破坏这些消费者可能导致面向客户的故障34。人们通常用数据契约来防止这类破坏。
将内部模式与外部模式解耦的一种常见方法是采用 发件箱模式(outbox pattern)。发件箱是具有独立模式的表;CDC 系统对外公开这些表,而不是数据库中的内部领域模型35、36。这样,开发者便可按需修改内部模式,而保持发件箱表不变。这看起来像双写——它确实就是双写。不过,两次写入都留在同一个系统(数据库)中,可以出现在同一个事务里,因此发件箱避开了 “保持系统同步” 中讨论的问题。
不过,发件箱也有一些权衡。开发者仍须维护内部模式与发件箱模式之间的转换,这可能并不容易。发件箱还会增加数据库写入底层存储的数据量,可能引发性能问题。
与变更数据捕获一样,重播事件日志可以重建系统当前状态。不过,两者处理日志压实的方式不同:
记录更新的 CDC 事件通常包含记录的完整新版本,因此一个主键的当前值完全由该主键的最新事件决定,日志压实可以丢弃同一主键的旧事件。
事件溯源的建模层次更高:事件通常表达用户操作的意图,而不是该操作引发状态更新的具体机制。后续事件通常不会覆盖先前事件,因此需要完整的事件历史才能重建最终状态,不能用同样的方式压实日志。
使用事件溯源的应用通常会保存由事件日志衍生出的当前状态快照,以免反复处理完整日志。不过,这只是一种性能优化,用来加快读取和崩溃恢复;系统的设计意图仍是永久保存所有原始事件,并能在需要时重新处理完整的事件日志。我们将在 “不变性的局限” 中讨论这一假设。
状态、流和不变性
我们在 第 11 章 中看到,批处理受益于输入文件的不变性:你可以在现有输入文件上运行实验性的处理作业,而不必担心损坏这些文件。正是不变性原则赋予了事件溯源和变更数据捕获如此强大的能力。
我们通常认为数据库存储着应用的当前状态——这种表示针对读取进行了优化,一般也最便于处理查询。状态的本质就在于它会变化,所以数据库除了插入数据,还支持更新和删除。这与不变性又如何相容?
只要状态会变化,它就是一段时间内各种事件改变它的结果。例如,当前可用座位列表取决于已经处理过的预订;当前账户余额取决于账户的贷记与借记;Web 服务器的响应时间图,则是所有已发生 Web 请求各自响应时间的聚合。
无论状态如何变化,总有一系列事件导致这些变化。事情可以做了再撤销,但那些事件确实发生过这一事实不会改变。关键在于,可变状态与不可变事件的仅追加日志并不矛盾,而是一枚硬币的两面。所有变化构成的日志——即 变更日志(changelog)——表示了状态如何随时间演化。
如果你偏爱数学,可以说应用状态是事件流对时间的积分,而变更流是状态对时间的微分,如 图 12-7 所示37、38。这个类比有其局限(例如,状态的二阶导数似乎没有什么意义),但它是思考数据的一个实用起点。

只要持久存储变更日志,状态便可以重现。如果把事件日志视为权威记录系统,把所有可变状态都看作由它衍生而来,系统中的数据流就更容易推理。正如 Jim Gray 和 Andreas Reuter 在 1992 年所说39:
从根本上说,根本没有必要保留数据库;日志已经包含了全部信息。之所以要存储数据库(即日志末尾所对应的当前状态),只是为了提升检索操作的性能。
日志压实是衔接日志与数据库状态的一种方式:它只保留每条记录的最新版本,丢弃已经被覆盖的版本。
不可变事件的优点
数据库中的不变性是一个古老的观念。例如,会计师几个世纪以来一直在财务记账中运用不变性。一笔交易发生后,会被记入仅追加的 分类账(ledger);分类账本质上就是事件日志,描述货币、商品或服务的转手。损益表、资产负债表等账目,则是汇总分类账中的交易而衍生出来的40。
如果出了差错,会计师不会删除或修改分类账中的错误交易,而是另加一笔交易来抵消错误,例如退还一笔误收的费用。错误交易会永远保留在分类账中,因为它对审计可能十分重要。如果根据错误分类账得出的错误数字已经公布,下一个会计期间的数字就会包含相应的更正。这在会计工作中再正常不过41。
这种可审计性对金融系统尤其重要,但许多不受严格监管的系统也能从中受益。如果你不慎部署了一段有缺陷的代码,把错误数据写进数据库,而这段代码还能以破坏性的方式覆盖数据,恢复起来会困难得多。有了不可变事件的仅追加日志,诊断事情经过并从问题中恢复就容易得多。同样,客服人员也可以利用审计日志诊断客户的请求和投诉。
不可变事件所包含的信息也比当前状态更加丰富。例如在购物网站上,顾客可能先把一件商品加入购物车,随后又将其移除。从履行订单的角度看,第二个事件抵消了第一个事件;但从分析角度看,知道顾客曾考虑购买某件商品、后来又放弃,也许很有用。或许他们日后会购买,或许他们找到了替代品。这些信息会留在事件日志中;如果数据库在商品移出购物车时就删除相应记录,信息也会随之丢失。
从同一事件日志中派生多个视图
此外,把可变状态与不可变事件日志分离之后,还可以从同一份事件日志衍生出几种面向不同读取方式的表示。这就像一股流有多个消费者一样(图 12-5):例如,分析数据库 Druid 会以这种方式直接从 Kafka 摄取数据,Kafka Connect 的汇聚连接器则可以把 Kafka 中的数据导出到各种数据库和索引33。
在事件日志与数据库之间加入显式的转换步骤,应用也更容易随时间演进。如果要引入一项新功能,用新的方式呈现现有数据,可以利用事件日志为新功能构建一个独立的、针对读取优化的视图,与现有系统并行运行,而不必修改现有系统。在许多情况下,新旧系统并行运行比在现有系统中执行复杂的模式迁移更加容易。等读取方全部切换到新系统、不再需要旧系统之后,只须关闭旧系统并回收资源即可42、43。
把数据写成一种针对写入优化的形式,再按需转换成多种针对读取优化的表示,这正是我们在 “事件溯源与 CQRS” 中见过的 命令查询职责分离(command query responsibility segregation,CQRS)模式。它不一定要求采用事件溯源:同样可以从 CDC 事件流构建多个物化视图44。
传统的数据库与模式设计方法建立在一个谬误之上:数据必须按将来查询它时所采用的形式写入。如果能把针对写入优化的事件日志转换成针对读取优化的应用状态,规范化与反规范化之争(参阅 “规范化、反规范化与连接”)就基本失去了意义。完全可以在读取优化视图中对数据做反规范化,因为转换过程提供了让视图与事件日志保持一致的机制。
在 “案例研究:社交网络首页时间线” 中,我们讨论过社交网络的主页时间线:它缓存着某位用户关注的人最近发布的帖子,就像邮箱一样。这也是针对读取优化的状态:主页时间线高度反规范化,因为你的帖子会复制到每位关注者的时间线中。不过,扇出服务会让这些重复状态与新帖子、新的关注关系保持同步,使这种重复仍然可控。
并发控制
CQRS 最大的缺点是事件日志的消费者通常以异步方式运行。因此,用户可能刚刚向日志写入数据,随即读取某个衍生视图,却发现该写入还没有反映到视图中。我们曾在 “读己之写” 中讨论过这个问题及其可能的解决办法。
一种解决办法是在把事件追加到日志时,同步更新读取视图。这要么需要在事件日志与衍生视图之间执行分布式事务,要么需要某种机制,等待事件反映到视图中。这两种方法通常都不切实际,因此视图一般还是异步更新。
另一方面,从事件日志衍生当前状态也简化了并发控制的某些方面。之所以经常需要多对象事务(参阅 “单对象与多对象操作”),是因为一次用户操作往往要修改多个不同位置的数据。采用事件溯源后,可以把一个事件设计成对用户操作的自包含描述。这样,用户操作只需在一个位置完成一次写入——把事件追加到日志——很容易保证其原子性。
如果事件日志和应用状态采用相同的分片方式(例如,处理分片 3 中某位客户的事件时,只需更新应用状态的分片 3),那么简单的单线程日志消费者无须对写入做任何并发控制——从设计上看,它一次只会处理一个事件(另见 “实际串行执行”)。日志在每个分片中定义了事件的串行顺序,从而消除了并发带来的不确定性27。如果一个事件涉及多个状态分片,就需要多做一些工作,我们将在 第 13 章 中讨论。
许多并未采用事件溯源模型的系统,同样依靠不变性来控制并发:各种数据库在内部使用不可变数据结构或多版本数据来支持时间点快照(参阅 “索引与快照隔离”)。Git、Mercurial、Fossil 等版本控制系统也依靠不可变数据保存文件的版本历史。
不变性的局限
永久保存所有变化的不可变历史,究竟在多大程度上可行?答案取决于数据集的变动量。有些工作负载以新增数据为主,很少更新或删除,很容易做成不可变的。另一些工作负载则在相对较小的数据集上频繁更新和删除;这时,不可变历史可能膨胀到难以承受,碎片化也可能成为问题,而压实和垃圾回收的性能会直接影响系统能否稳健运行45、46。
除了性能原因,有时还必须出于行政或法律原因删除数据,哪怕这有悖于不变性。例如,欧盟《通用数据保护条例》(GDPR)等隐私法规要求应用户请求删除其个人信息和错误信息;意外泄露敏感信息后,也可能需要控制影响范围。
在这些情况下,只在日志末尾追加一个事件,表示先前的数据应视为已删除,并不足够——你真正想做的是改写历史,假装那些数据从未写入。Datomic 把这种功能称为 切除(excision)47,Fossil 版本控制系统中也有一个类似概念,称为 排斥(shunning)48。
真正删除数据出乎意料地困难49,因为副本可能存在于许多地方。例如,存储引擎、文件系统和 SSD 往往把数据写到新位置,而不是在原地覆盖41;备份通常还会被刻意设计成不可变,以防意外删除或损坏。
一种允许删除不可变数据的方法是 密码学粉碎(crypto-shredding)50:把将来可能需要删除的数据加密存储;需要清除时,忘掉加密密钥。加密后的数据仍然存在,但已经无人能够使用。从某种意义上说,这只是转移了问题:实际数据现在不可变了,存储密钥的地方却是可变的。
此外,还必须事先决定哪些数据共用一把密钥,何时要改用不同密钥。这项决定非常重要,因为以后只能选择粉碎某把密钥加密的全部数据,或者一项也不粉碎,不能只删除其中一部分。如果为每个数据项分别保存一把密钥,密钥存储会变得与主数据存储一样庞大,难以管理。可穿刺加密(puncturable encryption)等更复杂的方案51可以选择性撤销一把密钥的部分解密能力,但尚未得到广泛应用。
总体而言,删除更像是“让数据更难取回”,而不是真正“让数据无法取回”。尽管如此,有时仍必须尝试,我们将在 “立法与自律” 中看到这一点。
流处理
到目前为止,本章已经讨论了流从何而来(用户活动事件、传感器和数据库写入),以及如何传输流(直接传递消息、通过消息代理传递,以及使用事件日志)。
接下来要讨论的是,拿到一股流之后能用它做什么——也就是如何处理它。大体上有三种选择:
取出事件中的数据,写入数据库、缓存、搜索索引或类似的存储系统,再供其他客户端查询。如 图 12-5 所示,这是一种让数据库与系统其他部分的变化保持同步的好办法,尤其是在流消费者是唯一写入数据库的客户端时。写入存储系统,相当于以流式方式完成 “批处理用例” 中讨论的工作。
以某种方式把事件推送给用户,例如发送告警邮件或推送通知,或者把事件流式传送到实时仪表板上加以可视化。这种情况下,人是流的最终消费者。
处理一股或多股输入流,产生一股或多股输出流。一股流可能先后经过由多个处理阶段组成的流水线,最终才到达某个输出(即选项 1 或 2)。
本章余下部分将讨论第三种选择:处理流并产生其他衍生流。执行这种流处理的代码称为 算子(operator)或 作业(job)。它与 第 11 章 中讨论的 Unix 进程和 MapReduce 作业关系密切,数据流模式也很相似:流处理器以只读方式消费输入流,再以仅追加方式把输出写到另一个位置。
流处理器中的分片和并行化模式,也与 第 11 章 介绍的 MapReduce 和数据流引擎非常相似,因此这里不再赘述。转换、过滤记录等基本映射操作的工作方式也相同。
流与批处理作业有一个关键区别:流永远不会结束。这个区别会带来许多后果。正如本章开头所说,对无界数据集进行排序没有意义,因此不能使用排序合并连接(参阅 “JOIN 与 GROUP BY”)。容错机制也必须改变:一个只运行了几分钟的批处理作业发生任务故障时,大可从头重启该任务;但一个流作业已经运行了几年,崩溃后再从头开始,通常不可行。
流处理的应用
长期以来,流处理一直用于监控:组织希望在特定事情发生时收到警报。例如:
欺诈检测系统需要判断信用卡的使用模式是否出现意外变化,并在信用卡可能被盗时将其冻结。
交易系统需要观察金融市场的价格变化,并按照指定规则执行交易。
制造系统需要监控工厂内机器的状态,一旦发生故障便迅速查明问题。
军事与情报系统需要追踪潜在侵略者的活动,发现攻击迹象时发出警报。
这类应用需要相当复杂的模式匹配与关联分析。不过,流处理也逐渐出现了其他用途。本节将简要比较其中几种应用。
复合事件处理
复合事件处理(complex event processing,CEP)是一种在 20 世纪 90 年代发展起来的事件流分析方法,特别适合需要搜索特定事件模式的应用52。正如正则表达式可以在字符串中搜索特定的字符模式,CEP 允许你指定规则,在流中搜索特定的事件模式。
CEP 系统通常使用 SQL 等高级声明式查询语言或图形用户界面,描述应该检测哪些事件模式。这些查询会提交给处理引擎;引擎消费输入流,并在内部维护一个状态机来执行所需的匹配。一旦找到匹配,引擎便发出一个 复合事件(complex event,名称由此而来),其中包含检测到的事件模式详情53。
这类系统中,查询与数据的关系恰好和普通数据库相反。数据库通常持久存储数据,把查询视为临时对象:查询到来时,数据库搜索与之匹配的数据,查询完成后便将其忘掉。CEP 引擎却反过来长期存储查询;每个事件到来时,引擎都会检查迄今所见的事件是否形成了与某个常驻查询相匹配的模式54。
CEP 的实现包括 Esper、Apama 和 TIBCO StreamBase。Flink、Spark Streaming 等分布式流处理器也支持使用 SQL 对流执行声明式查询。
流分析
流处理的另一个用途是对流进行 分析。CEP 与流分析的边界并不清晰,但一般来说,分析不太关心寻找特定的事件序列,而更关注大量事件上的聚合与统计指标,例如:
测量某类事件的速率(每个时间间隔发生多少次);
计算某段时间内一个值的滚动平均数;
将当前统计值与先前时间段对比(例如检测趋势,或在某项指标与上周同一时间相比异常偏高或偏低时发出警报)。
这类统计值通常在固定时间区间内计算。例如,你可能想知道过去 5 分钟内某项服务平均每秒收到多少次查询,以及这段时间内响应时间的第 99 百分位点。在几分钟内取平均,可以抹平相邻秒之间无关紧要的波动,同时仍能及时反映流量模式的变化。用于聚合的时间区间称为 窗口(window),我们将在 “时间推理” 中详细讨论。
流分析系统有时会使用概率算法,例如用布隆过滤器(我们在 “布隆过滤器” 中见过)判断集合成员关系,用 HyperLogLog55 估计基数,以及用各种算法估计百分位点(参阅 “计算百分位点”)。概率算法给出近似结果,但与精确算法相比,流处理器所需内存少得多。近似算法的这种用途有时让人误以为流处理系统总是有损而不精确,其实不然:流处理本身并没有任何近似性,使用概率算法只是一项优化56。
许多开源分布式流处理框架都以分析为设计目标,例如 Apache Storm、Spark Streaming、Flink、Samza、Apache Beam 和 Kafka Streams57。托管服务则包括 Google Cloud Dataflow 和 Azure Stream Analytics。
维护物化视图
我们已经看到,数据库的变更流可以用来维护缓存、搜索索引和数据仓库等衍生数据系统,使它们与源数据库保持同步。这些都是维护物化视图的例子:从某个数据集衍生出另一种视图,以便高效查询,并在底层数据变化时更新视图37。
同样,在事件溯源中,应用状态通过应用事件日志来维护;这里的应用状态也是一种物化视图。与流分析不同,只考虑某个时间窗口内的事件通常不够:除了日志压实可能丢弃的过时事件,构建物化视图可能需要任意时间段内的 所有 事件。实际上,你需要一个一直延伸到时间开端的窗口。
原则上,任何流处理器都可以用于维护物化视图。不过,有些面向分析的框架假定自己主要处理持续时间有限的窗口,而永久维护事件与这种假设背道而驰。Kafka Streams 和 Confluent 的 ksqlDB 建立在 Kafka 的日志压实支持之上,能够支持这类用途58。
数据库似乎很适合维护物化视图——毕竟,它们本来就是用来保存数据集完整副本的,而且许多数据库也支持物化视图。我们在 “物化视图与多维数据集” 中看到,数据仓库常见的分析查询可以物化成 OLAP 多维数据集。
遗憾的是,数据库通常通过批处理作业,或 PostgreSQL 的 REFRESH MATERIALIZED VIEW 之类的按需请求来刷新物化视图表。视图会定期重新计算,而不是在源数据更新时随之更新。这种方式有两个重大缺点,因而不适合通过流处理维护视图:
效率低下:每次更新视图都要重新处理所有数据,尽管绝大多数数据很可能没有变化。
数据不够新鲜:只有等到下一次计划更新重新运行查询,源数据的变化才会反映到物化视图中。
如果数据很容易分区,而且计算天然适合增量执行,也可以编写数据库触发器来高效更新物化视图。例如,如果物化视图维护每日销售总收入,那么每发生一笔新销售,只需更新相应日期的那一行。在少数场景中可以定制这样的解决方案,但许多 SQL 查询很难高效地转换为增量计算,甚至根本无法转换。
增量视图维护(incremental view maintenance,IVM)是解决上述问题的一种更通用方法。IVM 技术把 SQL 等关系语言转换成能够执行增量计算的算子。IVM 算法不再处理整个数据集,而只重新计算和更新发生变化的数据38、59、60。这样,视图计算的效率大幅提升,更新也可以更频繁地运行,显著改善数据新鲜度。
Materialize61、RisingWave、ClickHouse 和 Feldera 等数据库都采用 IVM 技术,提供高效的增量物化视图。这些数据库摄取事件流,实时提供物化视图。最近的事件缓存在内存中,并定期用于更新磁盘上的物化视图。读取时则把最近事件与已经物化的数据合并起来,提供一个统一的实时视图。由于读取通常用 SQL 表达,而物化视图往往以 OLAP 风格的格式存储,这些系统也支持 第 11 章 所讨论的大规模数据仓库式查询。
在流上搜索
CEP 可以搜索由多个事件构成的模式;除此以外,有时还需要按照全文搜索查询等复杂条件来搜索单个事件。
例如,媒体监测服务可以订阅媒体机构发布的新闻文章和节目源,搜索任何提及目标公司、产品或话题的新闻。为此,需要预先制定搜索查询,再不断让新闻条目流与该查询进行匹配。一些网站也有类似功能:例如,房地产网站的用户可以要求网站在市场上出现符合其搜索条件的新房源时通知他们。Elasticsearch 的 percolator 功能62 就是实现这类流式搜索的一种选择。
传统搜索引擎先为文档建立索引,再在索引上运行查询。搜索数据流却把这个过程颠倒过来:查询被存储下来,文档像 CEP 中的事件一样逐一流过这些查询。最简单的做法是让每份文档测试每条查询,但查询数量很大时会变慢。为了优化这个过程,也可以像索引文档那样索引查询,从而缩小可能匹配的查询集合63。
事件驱动架构与 RPC
在 “事件驱动的架构” 中,我们讨论过用消息传递系统代替 RPC,也就是把它用作服务之间的通信机制,Actor 模型便是一例。这些系统同样以消息和事件为基础,但我们通常不会把它们视为流处理器:
Actor 框架主要用于管理相互通信的模块如何并发和分布式执行,而流处理主要是一种数据管理技术。
Actor 之间的通信往往是短暂的一对一通信,而事件日志持久存在,并有多个订阅者。
Actor 可以采用任意方式通信,包括循环的请求/响应模式;流处理器通常组成无环流水线,每股流都是某个特定作业的输出,并从一组明确定义的输入流衍生而来。
不过,类 RPC 系统与流处理之间也有一些交叉。例如,Apache Storm 有一项称为 分布式 RPC 的功能,可以把用户查询分派给一组同时处理事件流的节点。这些查询会与输入流中的事件交错处理,结果再汇总并返回给用户(另见 “多分片数据处理”)。
Actor 框架也可以用来处理流。不过,许多这类框架无法保证发生崩溃时消息仍能送达;除非另外实现重试逻辑,否则处理过程不具备容错能力。
时间推理
流处理器经常要应对时间问题,尤其是用于分析时,往往会使用“过去五分钟的平均值”这样的时间窗口。“过去五分钟”听起来似乎清楚明确,实际上却出人意料地棘手。
批处理任务会迅速处理大量历史事件。如果需要按时间划分结果,批处理就必须查看每个事件中嵌入的时间戳。查看运行批处理的机器系统时钟毫无意义,因为任务的运行时间与事件的实际发生时间没有关系。
一个批处理任务可能几分钟内就读完一整年的历史事件;大多数情况下,我们关心的是这一年的历史时间线,而不是几分钟的处理时间。此外,使用事件中的时间戳还可以让处理具有确定性:针对同一输入重新运行同一个过程,会得到相同的结果。
另一方面,许多流处理框架使用处理机器的本地系统时钟(即 处理时间,processing time)来划分窗口64。这种方法简单明了;如果事件创建与事件处理之间的延迟短到可以忽略,也很合理。但只要处理延迟比较显著——也就是事件实际发生后过了一段明显可感的时间才得到处理——这种方法就会失效。
事件时间与处理时间
很多原因都会导致处理延迟:排队、网络故障、性能问题使消息代理或处理器发生争用、流消费者重启,以及故障恢复或修复代码缺陷后重新处理过去的事件。
消息延迟还会使消息以不可预测的顺序到达。例如,假设用户先发出一个 Web 请求,由 Web 服务器 A 处理;随后发出第二个请求,由服务器 B 处理。A 和 B 各自发出事件,描述自己处理的请求,但 B 的事件先于 A 的事件到达消息代理。于是,流处理器先看到 B 的事件,再看到 A 的事件,而它们实际发生的顺序恰好相反。
不妨拿《星球大战》系列电影来类比:第四部于 1977 年上映,第五部于 1980 年上映,第六部于 1983 年上映;随后依次是 1999、2002 和 2005 年上映的第一、二、三部,以及 2015、2017 和 2019 年上映的第七、八、九部65。如果按上映顺序观看,你处理这些电影的顺序就与故事的叙事顺序不同(集数好比事件时间戳,观看日期则是处理时间)。人类能够应对这种不连续性,但流处理算法必须专门设计,才能处理这类时间与顺序问题。
混淆 事件时间(event time)与处理时间会产生错误数据。例如,假设有一个流处理器用来测量请求速率(统计每秒请求数)。重新部署流处理器时,它可能停机一分钟,恢复运行后再处理积压的事件。如果按处理时间计算速率,处理积压期间看起来会突然出现异常的请求尖峰,而真实的请求速率其实一直很稳定(图 12-8)。

处理滞留事件
按事件时间定义窗口时,有一个棘手的问题:你永远无法确定某个窗口的所有事件是否已经到齐,还是仍有一些事件尚未到达。
例如,假设把事件分成一分钟的窗口,以统计每分钟的请求数。你已经统计了一批时间戳落在本小时第 37 分钟的事件;随着时间推移,新到事件现在大多落在第 38 和第 39 分钟。究竟什么时候才能宣布第 37 分钟的窗口已经结束,并输出计数器的值?
如果一段时间内没有再看到属于某个窗口的新事件,可以让它超时并宣布窗口就绪。然而,某些事件可能仍缓存在另一台机器上,因网络中断而延迟。你必须能够处理这些在窗口宣布完成后才到达的 滞留事件(straggler event)。大体上有两种选择1:
忽略滞留事件,因为正常情况下,它们可能只占所有事件的很小一部分。可以把丢弃事件的数量作为指标跟踪;如果开始丢弃大量数据,就发出警报。
发布 更正(correction):为窗口发布一个包含滞留事件的更新值。可能还需要撤回先前的输出。
有时可以用一条特殊消息表示:“从现在起,不会再有时间戳早于 t 的消息。”消费者可以利用它触发窗口66。然而,如果多台机器上的多个生产者都在生成事件,各自拥有不同的最小时间戳阈值,消费者就必须分别跟踪每个生产者。这时,增加或移除生产者会更加棘手。
你用的是谁的时钟?
如果事件会在系统中的多个位置缓冲,为事件赋予时间戳就更加困难。例如,考虑一个向服务器上报使用指标的移动应用。用户可能在设备离线时使用该应用;此时应用会把事件缓存在设备本地,等下次连上互联网时再发送给服务器,而这可能已是几小时甚至几天之后。对于流的消费者来说,这些事件看起来就像延迟极久的滞留事件。
这种情况下,事件时间戳其实应该是用户交互发生的时间,以移动设备的本地时钟为准。然而,用户控制的设备时钟往往不可信,因为它可能被无意或有意地设成错误时间(参阅 “时钟同步和准确性”)。服务器收到事件的时间以服务器时钟为准;由于服务器在你的控制之下,这个时间更可能准确,却不能很好地描述用户交互。
为了校正不准确的设备时钟,一种方法是记录三个时间戳67:
根据设备时钟,事件发生的时间;
根据设备时钟,事件发往服务器的时间;
根据服务器时钟,服务器收到事件的时间。
用第三个时间戳减去第二个,可以估算设备时钟与服务器时钟之间的偏移量(假设相对于所需的时间戳精度,网络延迟可以忽略)。再把这一偏移量应用到事件时间戳上,就能估算事件真正发生的时间——这里还要假设,从事件发生到事件发往服务器期间,设备时钟的偏移量没有变化。
这个问题并非流处理独有,批处理在时间推理上面临一模一样的问题。只是在流式环境中,我们更能意识到时间正在流逝,所以问题也更显眼。
窗口的类型
明确了如何确定事件时间戳之后,下一步就是决定如何定义时间窗口。窗口可用于聚合,例如统计事件数量,或计算窗口内各个值的平均数。以下几类窗口比较常见64、68:
- 滚动窗口(tumbling window)
滚动窗口的长度固定,每个事件恰好属于一个窗口。例如,使用一分钟的滚动窗口时,时间戳介于
10:03:00和10:03:59的所有事件归入一个窗口,介于10:04:00和10:04:59的事件归入下一个窗口,依此类推。实现一分钟滚动窗口时,可以把每个事件的时间戳向下取整到最近的整分钟,以确定它属于哪个窗口。- 跳跃窗口(hopping window)
跳跃窗口的长度也固定,但允许窗口重叠,以起到一定的平滑作用。例如,一个长度为五分钟、跳跃步长为一分钟的窗口,先包含
10:03:00到10:07:59之间的事件;下一个窗口包含10:04:00到10:08:59之间的事件,依此类推。可以先计算一分钟的滚动窗口,再聚合相邻的多个窗口,从而实现这个跳跃窗口。- 滑动窗口(sliding window)
滑动窗口包含彼此间隔不超过某段时间的所有事件。例如,五分钟的滑动窗口会同时覆盖发生在
10:03:39和10:08:12的事件,因为两者相差不到五分钟。注意,五分钟的滚动窗口或跳跃窗口使用固定边界,不一定会把这两个事件放在同一个窗口中。实现滑动窗口时,可以维护一个按时间排序的事件缓冲区,并在旧事件过期、离开窗口时将其移除。- 会话窗口(session window)
会话窗口与其他窗口不同,没有固定的持续时间。它把同一用户在时间上彼此接近的事件归为一组,并在用户有一段时间没有活动后结束窗口(例如,30 分钟内没有任何事件)。网站分析经常需要划分会话。
窗口操作通常需要维护临时状态。有些情况下,无论窗口多大、发生多少事件,状态的大小都是固定的:例如,计数操作不论窗口大小和事件数量如何,都只需要一个计数器。另一方面,滑动窗口以及下一节讨论的流连接,都必须缓冲事件,直到窗口结束。因此,窗口很大或流吞吐量很高时,流处理器可能需要保存大量临时状态。无论这些状态保存在内存还是磁盘上,都必须确保运行流处理任务的机器有足够容量来容纳它们。
流连接
在 “JOIN 与 GROUP BY” 中,我们讨论过批处理作业如何按键连接数据集,以及这种连接为何是数据流水线的重要组成部分。流处理把数据流水线推广到对无界数据集的增量处理,因此也同样需要对流执行连接。
不过,流中随时可能出现新事件,使流连接比批处理作业中的连接更具挑战。为了看清这个问题,我们把连接分成三类:流—流连接(stream-stream join)、流—表连接(stream-table join)和 表—表连接(table-table join)。下面各用一个例子来说明。
流流连接(窗口连接)
假设网站提供搜索功能,而你想发现最近的 URL 搜索趋势。每当有人输入搜索查询,就记录一个包含查询及返回结果的事件;每当有人点击某项搜索结果,又记录一个点击事件。为了计算搜索结果中每个 URL 的点击率,必须把搜索行为与点击行为的事件结合起来;它们可以通过相同的会话 ID 关联。广告系统也需要类似的分析69。
如果用户放弃搜索,点击也许永远不会发生;即使发生,搜索与点击之间的间隔也可能相差悬殊:通常只有几秒,但也可能长达数天或数周——例如用户搜索之后忘记了这个浏览器标签页,过了很久才回来点击某个结果。网络延迟不一,甚至可能让点击事件先于搜索事件到达。你可以为连接选择适当的窗口,例如只连接相隔不超过一小时的搜索与点击。
请注意,把搜索详情嵌入点击事件并不等同于连接两类事件:这样只能了解用户点击搜索结果的情况,却无法了解用户没有点击任何结果的搜索。衡量搜索质量需要准确的点击率,因此搜索事件和点击事件缺一不可。
为了实现这类连接,流处理器需要维护 状态(state),例如按会话 ID 索引过去一小时内的所有事件。每当搜索事件或点击事件到来,就把它加入相应索引,同时检查另一个索引,看看相同会话 ID 的另一事件是否已经到达。找到匹配时,发出一个事件,说明哪项搜索结果被点击;如果搜索事件过期时仍未看到匹配的点击事件,则发出一个事件,说明哪些搜索结果没有被点击。
流表连接(流扩充)
在 “JOIN 与 GROUP BY”(图 11-2)中,我们见过批处理作业连接两个数据集的例子:一组用户活动事件和一个用户档案数据库。很自然地,可以把用户活动事件视为一股流,在流处理器中持续执行同样的连接:输入是包含用户 ID 的活动事件流,输出则是活动事件流,其中的用户 ID 已经补充了相应的用户档案信息。这个过程有时称为用数据库中的信息 扩充(enrich)活动事件。
执行这项连接时,流处理器要逐一查看活动事件,在数据库中查找事件里的用户 ID,再把档案信息加入活动事件。数据库查找可以通过查询远程数据库来实现;不过,正如 “JOIN 与 GROUP BY” 中所讨论的,这类远程查询很可能速度缓慢,还可能使数据库过载58。
另一种做法是把数据库副本载入流处理器,在本地查询,免去网络往返。由于数据库的本地副本可能是内存散列表(如果足够小),也可能是本地磁盘上的索引,因此这项技术称为 散列连接(hash join)。
它与批处理作业的区别在于:批处理作业把数据库某个时间点的快照用作输入;流处理器却长期运行,而数据库内容很可能随时间变化,所以流处理器中的本地副本必须持续更新。变更数据捕获可以解决这个问题:除了活动事件流,流处理器还可以订阅用户档案数据库的变更日志。每当创建或修改档案时,流处理器就更新本地副本。这样,我们实际上得到了两股流之间的连接:活动事件与档案更新。
流—表连接其实与流—流连接很相似。最大的区别在于,对表的变更日志流执行连接时,使用的是一个回溯至“时间开端”的窗口(概念上是无限窗口),其中记录的新版本会覆盖旧版本;对另一股输入流,连接可能根本不维护窗口。
表表连接(维护物化视图)
考虑 “案例研究:社交网络首页时间线” 中讨论的社交网络时间线。我们说过,用户查看主页时间线时,如果遍历他所关注的所有人,找出他们最近发布的帖子再合并,代价实在太高。
我们需要的是时间线缓存:为每位用户准备一个“收件箱”,帖子发布时便写入其中,这样读取时间线只需查找一次。物化并维护这个缓存,需要处理以下事件:
用户 u 发布新帖子时,把帖子加入所有关注 u 的用户的时间线。
用户删除一篇帖子或删除整个账户时,从所有用户的时间线中移除相应帖子。
用户 u
1开始关注用户 u2时,把 u2最近的帖子加入 u1的时间线。用户 u
1取消关注用户 u2时,从 u1的时间线中移除 u2的帖子。
要在流处理器中维护这项缓存,需要一股帖子事件流(发布和删除),以及一股关注关系事件流(关注和取消关注)。流处理器还要维护一个数据库,记录每位用户的关注者集合,以便新帖子到来时知道应该更新哪些时间线。
也可以换个角度来看:这项流处理维护着一个连接两张表(帖子和关注关系)的查询物化视图,大致如下:
流之间的连接直接对应查询中的表连接。时间线实际上就是查询结果的缓存,每当底层表发生变化时都会更新。
连接的时间依赖性
这里介绍的三类连接(流—流、流—表和表—表)有许多共同点:流处理器都要根据连接的一侧维护某种状态(搜索和点击事件、用户档案或关注者列表),再在连接另一侧的消息到来时查询该状态。
维护状态的事件顺序十分重要——先关注再取消关注,与顺序相反的结果不同。在 Kafka 这样的分片事件日志中,同一个分片(分区)内的事件顺序能够保持,但不同流或不同分片之间通常没有顺序保证。
这就带来一个问题:如果不同流上的事件发生时间相近,应该按什么顺序处理?在流—表连接的例子中,如果用户更新了档案,哪些活动事件应该与旧档案连接(在档案更新前处理),哪些应该与新档案连接(在档案更新后处理)?换句话说,如果状态随时间变化,而你要与这个状态连接,究竟应该使用哪个时间点的状态?
这种时间依赖会出现在许多地方。例如,销售商品时要为发票采用正确的税率;税率取决于国家或州、产品类型和销售日期,因为税率会不时变化。把销售记录与税率表连接时,通常需要采用销售发生时的税率;如果正在重新处理历史数据,它可能与当前税率不同。
如果不同流之间的事件顺序不确定,连接也会变得不确定70。这意味着,即使针对同一输入重新运行同一个作业,也不一定得到相同结果:再次运行时,各输入流上的事件可能以不同方式交错。
在数据仓库中,这个问题称为 缓慢变化维度(slowly changing dimension,SCD),通常为所连接记录的每个具体版本赋予唯一标识符来解决。例如,每次税率变化时,都给新税率分配一个新标识符;发票则包含销售时所用税率的标识符71、72。这样连接就具有确定性,但也无法再做日志压实,因为必须保留表中记录的所有版本。另一种方法是反规范化数据,把适用税率直接放入每个销售事件中。
容错
本章最后一节来看看流处理器如何容忍故障。我们在 第 11 章 中看到,批处理框架相当容易实现容错:一项任务失败后,只须在另一台机器上重新启动,并丢弃失败任务的输出。之所以能透明地重试,是因为输入文件不可变,每项任务都把输出写入独立文件,而且只有任务成功完成后,输出才会变得可见。
具体来说,批处理的容错方式可以确保:即使某些任务实际上失败过,批处理作业的输出也与没有发生过任何问题时相同。看起来每条输入记录都只处理了恰好一次——没有记录被跳过,也没有记录被处理两遍。重新启动任务意味着记录实际上可能处理了多次,但输出中可见的效果却像只处理了一次。这个原则称为 恰好一次语义(exactly-once semantics),不过 等效一次(effectively-once)或许是更贴切的叫法73。
流处理也面临同样的容错问题,却没那么容易解决:不能等任务完成后才让输出可见,因为流是无限的,永远也处理不完。
微批处理与检查点
一种解决办法是把流分成小块,把每一块当作一个微型批处理。这种方法称为 微批处理(microbatching),Spark Streaming 就采用了它74。批次通常约为一秒,这是性能权衡的结果:批次越小,调度与协调开销越大;批次越大,流处理器的结果变得可见之前,延迟就越长。
微批处理还隐式提供了一个大小等于批次大小的滚动窗口——它按处理时间而非事件时间戳划分。需要更大窗口的作业,必须显式地把状态从一个微批次传递到下一个微批次。
Apache Flink 采用一种变体:定期生成滚动的状态检查点,并将其写入持久存储75、76。如果流算子崩溃,可以从最近的检查点重新启动,并丢弃从上一个检查点到崩溃之间产生的所有输出。检查点由消息流中的屏障触发,类似于微批次之间的边界,但不会强制规定窗口大小。
在流处理框架内部,微批处理与检查点都能提供与批处理相同的恰好一次语义。然而,一旦输出离开流处理器——例如写入数据库、向外部消息代理发送消息或发送电子邮件——框架便无法丢弃失败微批次的输出。这时,重启失败任务会使外部副作用发生两次,仅靠微批处理或检查点不足以避免这个问题。
再谈原子提交
为了在发生故障时营造恰好一次处理的效果,必须确保处理事件产生的所有输出和副作用,当且仅当 处理成功时才生效。这些影响包括:发给下游算子或外部消息传递系统的消息(包括电子邮件或推送通知)、数据库写入、算子状态的变化,以及对输入消息的确认应答(包括推进基于日志的消息代理中的消费者偏移量)。
这些事情要么全部以原子方式发生,要么一件也不发生,不能彼此失去同步。这种做法听起来似曾相识,是因为我们曾在分布式事务和两阶段提交的语境下讨论过它(参阅 “恰好一次消息处理”)。
我们在 第 10 章 中讨论了 XA 等传统分布式事务实现的问题。不过,在限制更严格的环境中,可以高效实现这样的原子提交机制。Google Cloud Dataflow66、75、VoltDB77 和 Apache Kafka78、79 都采用了这种方法。与 XA 不同,这些实现不会试图跨异构技术提供事务,而是由流处理框架同时管理状态变化与消息传递,把事务留在框架内部。还可以在单个事务中处理多条输入消息,从而分摊事务协议的开销。
幂等性
我们的目标是丢弃所有失败任务的部分输出,以便能安全地重试,而不会生效两次。分布式事务是实现这个目标的一种方法;另一种方法是依靠 幂等性(idempotence),正如 “持久化执行与工作流” 中所见80。
幂等操作可以执行多次,效果与只执行一次相同。例如,删除键值存储中的一个键是幂等的(再次删除不会产生进一步影响);递增计数器却不是幂等的(再执行一次递增,值就增加了两次)。
即使操作本身并非天然幂等,通常也可以用一点额外元数据让它变得幂等。例如,消费 Kafka 消息时,每条消息都有一个持久、单调递增的偏移量。向外部数据库写入值时,可以把触发最近一次写入的消息偏移量与值一并写入。这样便能判断某项更新是否已经应用,避免再次执行同一更新。
Storm 的 Trident 也用类似思路处理状态。依靠幂等性包含几个假设:重启失败任务时,必须以相同顺序重播相同消息(基于日志的消息代理可以做到);处理过程必须具有确定性;而且不能有其他节点并发更新同一个值81、82。
从一个处理节点故障切换到另一个节点时,可能还需要采用栅栏机制(参阅 “分布式锁和租约”),以防某个被认为已经死亡、实际仍然存活的节点造成干扰。尽管有这么多限制,幂等操作仍是实现恰好一次语义的有效方式,而且开销很小。
失败后重建状态
任何需要状态的流处理过程——例如计数器、平均值、直方图等窗口聚合,以及连接所用的表和索引——都必须确保发生故障后能够恢复状态。
一种选择是把状态保存在远程数据存储中并进行复制,不过为每条消息查询一次远程数据库可能很慢。另一种选择是把状态保存在流处理器本地,再定期复制。这样,流处理器从故障中恢复时,新任务可以读取状态副本,继续处理而不丢失数据。
例如,Flink 定期捕获算子状态快照,将其写入分布式文件系统等持久存储75、76;Kafka Streams 则把状态变化发送到一个启用了日志压实的专用 Kafka 主题,以类似变更数据捕获的方式复制状态83。VoltDB 会在多个节点上冗余处理每条输入消息,以此复制状态(参阅 “实际串行执行”)。
有些情况下,甚至不必复制状态,因为可以根据输入流重建。例如,如果状态只是较短窗口上的聚合,重播该窗口对应的输入事件也许足够快。如果状态是由变更数据捕获维护的数据库本地副本,也可以从经过日志压实的变更流重建数据库。
不过,所有这些权衡都取决于底层基础设施的性能特征:有些系统的网络延迟可能低于磁盘访问延迟,网络带宽也可能与磁盘带宽相当。不存在适用于所有情况的理想权衡;随着存储和网络技术演进,本地状态与远程状态各自的优势也可能发生变化。
本章小结
本章讨论了事件流、它们的用途,以及如何处理它们。从某种意义上说,流处理与 第 11 章 讨论的批处理非常相似,只不过流处理持续作用于无界(永不终止)的流,而不是大小固定的输入84。从这个角度看,消息代理和事件日志就是文件系统的流式对应物。
我们花了一些篇幅比较两类消息代理:
- AMQP/JMS 风格的消息代理
代理把每条消息分别分配给消费者;消费者成功处理一条消息后,就对它进行确认应答。消息得到确认后,代理会将其删除。这种方法适合用作异步 RPC(另见 “事件驱动的架构”),例如用于任务队列:消息处理的确切顺序并不重要,处理完成后也无须返回重读旧消息。
- 基于日志的消息代理
代理把一个分片中的所有消息分配给同一个消费者节点,并始终按相同顺序传递消息。系统通过分片实现并行;消费者则为已处理的最后一条消息偏移量创建检查点,以此跟踪进度。代理把消息保留在磁盘上,因此必要时可以跳回去重新读取旧消息。
基于日志的方法与数据库的复制日志(参阅 第 6 章)和日志结构存储引擎(参阅 第 4 章)有相似之处。正如 第 10 章 所述,它也是共识的一种形式。我们看到,这种方法尤其适合那些消费输入流,并生成衍生状态或衍生输出流的流处理系统。
至于流从何而来,我们讨论了几种可能:用户活动事件、定期提供读数的传感器,以及数据源(例如金融市场数据),都天然适合表示成流。把数据库写入视为流同样很有用:可以通过变更数据捕获隐式取得变更日志——即数据库所有变化的历史——也可以通过事件溯源显式取得。日志压实让流可以保留数据库内容的完整副本。
把数据库表示成流,为系统集成带来了强大能力。消费变更日志并把变化应用到衍生系统,就能让搜索索引、缓存、分析系统等衍生数据系统持续保持最新。甚至可以从头开始消费变更日志,一直追到当前时刻,从而基于现有数据构建全新的视图。
以流的形式维护状态、重播消息的能力,也是各种流处理框架实现流连接和容错技术的基础。我们讨论了流处理的几种用途,包括搜索事件模式(复合事件处理)、计算窗口聚合(流分析),以及让衍生数据系统保持最新(物化视图)。
随后,我们讨论了流处理器进行时间推理时的困难,包括处理时间与事件时间戳的区别,以及如何处理那些在窗口本以为已经完成后才姗姗来迟的滞留事件。
我们区分了流处理过程中可能出现的三类连接:
- 流—流连接
两股输入流都由活动事件组成,连接算子会在某个时间窗口内寻找彼此相关的事件。例如,它可以匹配同一用户在 30 分钟内执行的两个操作。如果想寻找同一股流中的相关事件,连接的两侧实际上也可以是同一股流,这称为 自连接(self-join)。
- 流—表连接
一股输入流由活动事件组成,另一股则是数据库变更日志。变更日志用来让数据库的本地副本保持最新。对于每个活动事件,连接算子查询数据库,并输出经过扩充的活动事件。
- 表—表连接
两股输入流都是数据库变更日志。这种情况下,任意一侧的每次变化都与另一侧的最新状态连接。结果是两张表连接所得物化视图的变更流。
最后,我们讨论了流处理器实现容错和恰好一次语义的技术。与批处理一样,必须丢弃失败任务产生的部分输出。但流处理过程长期运行并持续产生输出,无法简单地丢弃所有输出。因此,需要微批处理、检查点、事务或幂等写入等机制,在更细的粒度上进行恢复。
脚注
参考文献
Tyler Akidau, Robert Bradshaw, Craig Chambers, Slava Chernyak, Rafael J. Fernández-Moctezuma, Reuven Lax, Sam McVeety, Daniel Mills, Frances Perry, Eric Schmidt, and Sam Whittle. The Dataflow Model: A Practical Approach to Balancing Correctness, Latency, and Cost in Massive-Scale, Unbounded, Out-of-Order Data Processing. Proceedings of the VLDB Endowment, volume 8, issue 12, pages 1792–1803, August 2015. doi:10.14778/2824032.2824076 ↩︎ ↩︎
Harold Abelson, Gerald Jay Sussman, and Julie Sussman. Structure and Interpretation of Computer Programs, 2nd edition. MIT Press, 1996. ISBN: 978-0-262-51087-5, archived at archive.org/details/sicp_20211010 ↩︎
Patrick Th. Eugster, Pascal A. Felber, Rachid Guerraoui, and Anne-Marie Kermarrec. The Many Faces of Publish/Subscribe. ACM Computing Surveys, volume 35, issue 2, pages 114–131, June 2003. doi:10.1145/857076.857078 ↩︎
Don Carney, Uğur Çetintemel, Mitch Cherniack, Christian Convey, Sangdon Lee, Greg Seidman, Michael Stonebraker, Nesime Tatbul, and Stan Zdonik. Monitoring Streams – A New Class of Data Management Applications. At 28th International Conference on Very Large Data Bases (VLDB), August 2002. doi:10.1016/B978-155860869-6/50027-5 ↩︎
Matthew Sackman. Pushing Back. wellquite.org, May 2016. Archived at perma.cc/3KCZ-RUFY ↩︎ ↩︎
Thomas Figg (tef). how (not) to write a pipeline. cohost.org, June 2023. Archived at perma.cc/A3V8-NYCM ↩︎
Vicent Martí. Brubeck, a statsd-Compatible Metrics Aggregator. github.blog, June 2015. Archived at perma.cc/TP3Q-DJYM ↩︎
Seth Lowenberger. MoldUDP64 Protocol Specification V 1.00. nasdaqtrader.com, July 2009. Archived at https://perma.cc/7CRQ-QBD7 ↩︎
Ian Malpass. Measure Anything, Measure Everything. codeascraft.com, February 2011. Archived at archive.org ↩︎
Dieter Plaetinck. 25 Graphite, Grafana and statsd Gotchas. grafana.com, March 2016. Archived at perma.cc/3NP3-67U7 ↩︎
Jeff Lindsay. Web Hooks to Revolutionize the Web. progrium.com, May 2007. Archived at perma.cc/BF9U-XNX4 ↩︎
Jim N. Gray. Queues Are Databases. Microsoft Research Technical Report MSR-TR-95-56, December 1995. Archived at arxiv.org ↩︎
Mark Hapner, Rich Burridge, Rahul Sharma, Joseph Fialli, Kate Stout, and Nigel Deakin. JSR-343 Java Message Service (JMS) 2.0 Specification. jms-spec.java.net, March 2013. Archived at perma.cc/E4YG-46TA ↩︎
Sanjay Aiyagari, Matthew Arrott, Mark Atwell, Jason Brome, Alan Conway, Robert Godfrey, Robert Greig, Pieter Hintjens, John O’Hara, Matthias Radestock, Alexis Richardson, Martin Ritchie, Shahrokh Sadjadi, Rafael Schloming, Steven Shaw, Martin Sustrik, Carl Trieloff, Kim van der Riet, and Steve Vinoski. AMQP: Advanced Message Queuing Protocol Specification. Version 0-9-1, November 2008. Archived at perma.cc/6YJJ-GM9X ↩︎
Architectural overview of Pub/Sub. cloud.google.com, 2025. Archived at perma.cc/VWF5-ABP4 ↩︎ ↩︎
Aris Tzoumas. Lessons from scaling PostgreSQL queues to 100k events per second. rudderstack.com, July 2025. Archived at perma.cc/QD8C-VA4Y ↩︎
Robin Moffatt. Kafka Connect Deep Dive – Error Handling and Dead Letter Queues. confluent.io, March 2019. Archived at perma.cc/KQ5A-AB28 ↩︎
Dunith Danushka. Message reprocessing: How to implement the dead letter queue. redpanda.com. Archived at perma.cc/R7UB-WEWF ↩︎
Damien Gasparina, Loic Greffier, and Sebastien Viale. KIP-1034: Dead letter queue in Kafka Streams. cwiki.apache.org, April 2024. Archived at perma.cc/3VXV-QXAN ↩︎
Jay Kreps, Neha Narkhede, and Jun Rao. Kafka: A Distributed Messaging System for Log Processing. At 6th International Workshop on Networking Meets Databases (NetDB), June 2011. Archived at perma.cc/CSW7-TCQ5 ↩︎
Jay Kreps. Benchmarking Apache Kafka: 2 Million Writes Per Second (On Three Cheap Machines). engineering.linkedin.com, April 2014. Archived at archive.org ↩︎
Kartik Paramasivam. How We’re Improving and Advancing Kafka at LinkedIn. engineering.linkedin.com, September 2015. Archived at perma.cc/3S3V-JCYJ ↩︎
Philippe Dobbelaere and Kyumars Sheykh Esmaili. Kafka versus RabbitMQ: A comparative study of two industry reference publish/subscribe implementations. At 11th ACM International Conference on Distributed and Event-based Systems (DEBS), June 2017. doi:10.1145/3093742.3093908 ↩︎
Kate Holterhoff. Why Message Queues Endure: A History. redmonk.com, December 2024. Archived at perma.cc/6DX8-XK4W ↩︎
Andrew Schofield. KIP-932: Queues for Kafka. cwiki.apache.org, May 2023. Archived at perma.cc/LBE4-BEMK ↩︎
Jack Vanlightly. The advantages of queues on logs. jack-vanlightly.com, October 2023. Archived at perma.cc/WJ7V-287K ↩︎
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 ↩︎ ↩︎
Andy Hattemer. Change Data Capture is having a moment. Why? materialize.com, September 2021. Archived at perma.cc/AL37-P53C ↩︎
Prem Santosh Udaya Shankar. Streaming MySQL Tables in Real-Time to Kafka. engineeringblog.yelp.com, August 2016. Archived at perma.cc/5ZR3-2GVV ↩︎
Andreas Andreakis, Ioannis Papapanagiotou. DBLog: A Watermark Based Change-Data-Capture Framework. October 2020. Archived at arxiv.org ↩︎
Jiri Pechanec. Percolator. debezium.io, October 2021. Archived at perma.cc/EQ8E-W6KQ ↩︎
Debezium maintainers. Debezium Connector for Cassandra. debezium.io. Archived at perma.cc/WR6K-EKMD ↩︎
Neha Narkhede. Announcing Kafka Connect: Building Large-Scale Low-Latency Data Pipelines. confluent.io, February 2016. Archived at perma.cc/8WXJ-L6GF ↩︎ ↩︎
Chris Riccomini. Kafka change data capture breaks database encapsulation. cnr.sh, November 2018. Archived at perma.cc/P572-9MKF ↩︎
Gunnar Morling. “Change Data Capture Breaks Encapsulation”. Does it, though? decodable.co, November 2023. Archived at perma.cc/YX2P-WNWR ↩︎
Gunnar Morling. Revisiting the Outbox Pattern. decodable.co, October 2024. Archived at perma.cc/M5ZL-RPS9 ↩︎
Ashish Gupta and Inderpal Singh Mumick. Maintenance of Materialized Views: Problems, Techniques, and Applications. IEEE Data Engineering Bulletin, volume 18, issue 2, pages 3–18, June 1995. Archived at archive.org ↩︎ ↩︎ ↩︎
Mihai Budiu, Tej Chajed, Frank McSherry, Leonid Ryzhyk, Val Tannen. DBSP: Incremental Computation on Streams and Its Applications to Databases. SIGMOD Record, volume 53, issue 1, pages 87–95, March 2024. doi:10.1145/3665252.3665271 ↩︎ ↩︎
Jim Gray and Andreas Reuter. Transaction Processing: Concepts and Techniques. Morgan Kaufmann, 1992. ISBN: 9781558601901 ↩︎
Martin Kleppmann. Accounting for Computer Scientists. martin.kleppmann.com, March 2011. Archived at perma.cc/9EGX-P38N ↩︎
Pat Helland. Immutability Changes Everything. At 7th Biennial Conference on Innovative Data Systems Research (CIDR), January 2015. ↩︎ ↩︎
Martin Kleppmann. Making Sense of Stream Processing. Report, O’Reilly Media, May 2016. Archived at perma.cc/RAY4-JDVX ↩︎
Kartik Paramasivam. Stream Processing Hard Problems – Part 1: Killing Lambda. engineering.linkedin.com, June 2016. Archived at archive.org ↩︎
Stéphane Derosiaux. CQRS: What? Why? How? sderosiaux.medium.com, September 2019. Archived at perma.cc/FZ3U-HVJ4 ↩︎
Baron Schwartz. Immutability, MVCC, and Garbage Collection. xaprb.com, December 2013. Archived at archive.org ↩︎
Daniel Eloff, Slava Akhmechet, Jay Kreps, et al. Re: Turning the Database Inside-out with Apache Samza. Hacker News discussion, news.ycombinator.com, March 2015. Archived at perma.cc/ML9E-JC83 ↩︎
Datomic Documentation: Excision. Cognitect, Inc., docs.datomic.com. Archived at perma.cc/J5QQ-SH32 ↩︎
Fossil Documentation: Deleting Content from Fossil. fossil-scm.org, 2025. Archived at perma.cc/DS23-GTNG ↩︎
Jay Kreps. The irony of distributed systems is that data loss is really easy but deleting data is surprisingly hard. x.com, March 2015. Archived at perma.cc/7RRZ-V7B7 ↩︎
Brent Robinson. Crypto shredding: How it can solve modern data retention challenges. medium.com, January 2019. Archived at https://perma.cc/4LFK-S6XE ↩︎
Matthew D. Green and Ian Miers. Forward Secure Asynchronous Messaging from Puncturable Encryption. At IEEE Symposium on Security and Privacy, May 2015. doi:10.1109/SP.2015.26 ↩︎
David C. Luckham. What’s the Difference Between ESP and CEP? complexevents.com, June 2019. Archived at perma.cc/E7PZ-FDEF ↩︎
Arvind Arasu, Shivnath Babu, and Jennifer Widom. The CQL Continuous Query Language: Semantic Foundations and Query Execution. The VLDB Journal, volume 15, issue 2, pages 121–142, June 2006. doi:10.1007/s00778-004-0147-z ↩︎
Julian Hyde. Data in Flight: How Streaming SQL Technology Can Help Solve the Web 2.0 Data Crunch. ACM Queue, volume 7, issue 11, December 2009. doi:10.1145/1661785.1667562 ↩︎
Philippe Flajolet, Éric Fusy, Olivier Gandouet, and Frédéric Meunier. HyperLogLog: The Analysis of a Near-Optimal Cardinality Estimation Algorithm. At Conference on Analysis of Algorithms (AofA), June 2007. doi:10.46298/dmtcs.3545 ↩︎
Jay Kreps. Questioning the Lambda Architecture. oreilly.com, July 2014. Archived at perma.cc/2WY5-HC8Y ↩︎
Ian Reppel. An Overview of Apache Streaming Technologies. ianreppel.org, March 2016. Archived at perma.cc/BB3E-QJLW ↩︎
Jay Kreps. Why Local State is a Fundamental Primitive in Stream Processing. oreilly.com, July 2014. Archived at perma.cc/P8HU-R5LA ↩︎ ↩︎
RisingWave Labs. Deep Dive Into the RisingWave Stream Processing Engine - Part 2: Computational Model. risingwave.com, November 2023. Archived at perma.cc/LM74-XDEL ↩︎
Frank McSherry, Derek G. Murray, Rebecca Isaacs, and Michael Isard. Differential dataflow. At 6th Biennial Conference on Innovative Data Systems Research (CIDR), January 2013. ↩︎
Andy Hattemer. Incremental Computation in the Database. materialize.com, March 2020. Archived at perma.cc/AL94-YVRN ↩︎
Shay Banon. Percolator. elastic.co, February 2011. Archived at perma.cc/LS5R-4FQX ↩︎
Alan Woodward and Martin Kleppmann. Real-Time Full-Text Search with Luwak and Samza. martin.kleppmann.com, April 2015. Archived at perma.cc/2U92-Q7R4 ↩︎
Tyler Akidau. The World Beyond Batch: Streaming 102. oreilly.com, January 2016. Archived at perma.cc/4XF9-8M2K ↩︎ ↩︎
Stephan Ewen. Streaming Analytics with Apache Flink. At Kafka Summit, April 2016. Archived at perma.cc/QBQ4-F9MR ↩︎
Tyler Akidau, Alex Balikov, Kaya Bekiroğlu, Slava Chernyak, Josh Haberman, Reuven Lax, Sam McVeety, Daniel Mills, Paul Nordstrom, and Sam Whittle. MillWheel: Fault-Tolerant Stream Processing at Internet Scale. Proceedings of the VLDB Endowment, volume 6, issue 11, pages 1033–1044, August 2013. doi:10.14778/2536222.2536229 ↩︎ ↩︎
Alex Dean. Improving Snowplow’s Understanding of Time. snowplow.io, September 2015. Archived at perma.cc/6CT9-Z3Q2 ↩︎
Azure Stream Analytics: Windowing functions. Microsoft Azure Reference, learn.microsoft.com, July 2025. Archived at archive.org ↩︎
Rajagopal Ananthanarayanan, Venkatesh Basker, Sumit Das, Ashish Gupta, Haifeng Jiang, Tianhao Qiu, Alexey Reznichenko, Deomid Ryabkov, Manpreet Singh, and Shivakumar Venkataraman. Photon: Fault-Tolerant and Scalable Joining of Continuous Data Streams. At ACM International Conference on Management of Data (SIGMOD), June 2013. doi:10.1145/2463676.2465272 ↩︎
Ben Kirwin. Doing the Impossible: Exactly-Once Messaging Patterns in Kafka. ben.kirw.in, November 2014. Archived at perma.cc/A5QL-QRX7 ↩︎
Pat Helland. Data on the Outside Versus Data on the Inside. At 2nd Biennial Conference on Innovative Data Systems Research (CIDR), January 2005. ↩︎
Ralph Kimball and Margy Ross. The Data Warehouse Toolkit: The Definitive Guide to Dimensional Modeling, 3rd edition. John Wiley & Sons, 2013. ISBN: 978-1-118-53080-1 ↩︎
Viktor Klang. I’m coining the phrase ’effectively-once’ for message processing with at-least-once + idempotent operations. x.com, October 2016. Archived at perma.cc/7DT9-TDG2 ↩︎
Matei Zaharia, Tathagata Das, Haoyuan Li, Scott Shenker, and Ion Stoica. Discretized Streams: An Efficient and Fault-Tolerant Model for Stream Processing on Large Clusters. At 4th USENIX Conference in Hot Topics in Cloud Computing (HotCloud), June 2012. ↩︎
Kostas Tzoumas, Stephan Ewen, and Robert Metzger. High-Throughput, Low-Latency, and Exactly-Once Stream Processing with Apache Flink. ververica.com, August 2015. Archived at archive.org ↩︎ ↩︎ ↩︎
Paris Carbone, Gyula Fóra, Stephan Ewen, Seif Haridi, and Kostas Tzoumas. Lightweight Asynchronous Snapshots for Distributed Dataflows. arXiv:1506.08603 [cs.DC], June 2015. ↩︎ ↩︎
Ryan Betts and John Hugg. Fast Data: Smart and at Scale. Report, O’Reilly Media, October 2015. Archived at perma.cc/VQ6S-XQQY ↩︎
Neha Narkhede and Guozhang Wang. Exactly-Once Semantics Are Possible: Here’s How Kafka Does It. confluent.io, June 2019. Archived at perma.cc/Q2AU-Q2ED ↩︎
Jason Gustafson, Flavio Junqueira, Apurva Mehta, Sriram Subramanian, and Guozhang Wang. KIP-98 – Exactly Once Delivery and Transactional Messaging. cwiki.apache.org, November 2016. Archived at perma.cc/95PT-RCTG ↩︎
Pat Helland. Idempotence Is Not a Medical Condition. Communications of the ACM, volume 55, issue 5, page 56, May 2012. doi:10.1145/2160718.2160734 ↩︎
Jay Kreps. Re: Trying to Achieve Deterministic Behavior on Recovery/Rewind. Email to samza-dev mailing list, September 2014. Archived at perma.cc/7DPD-GJNL ↩︎
E. N. (Mootaz) Elnozahy, Lorenzo Alvisi, Yi-Min Wang, and David B. Johnson. A Survey of Rollback-Recovery Protocols in Message-Passing Systems. ACM Computing Surveys, volume 34, issue 3, pages 375–408, September 2002. doi:10.1145/568522.568525 ↩︎
Adam Warski. Kafka Streams – How Does It Fit the Stream Processing Landscape? softwaremill.com, June 2016. Archived at perma.cc/WQ5Q-H2J2 ↩︎
Stephan Ewen, Fabian Hueske, and Xiaowei Jiang. Batch as a Special Case of Streaming and Alibaba’s contribution of Blink. flink.apache.org, February 2019. Archived at perma.cc/A529-SKA9 ↩︎