跳转到内容
← 返回核心概念
系统与架构计算机科学 · 分布式系统16 分钟阅读

消息队列与流处理

Message Queues and Streaming

2010 年,LinkedIn 的工程团队遇到了一个奇怪的问题:系统里有数十个数据管道,把数据从各个来源同步到各个目的地,形成了一张复杂的点对点网络。每增加一个新数据来源或新的数据消费者,就需要新写一批同步管道,整个系统呈现出 $O(n^2)$ 的连接复杂度。 他们的解决方案是发明一个"统一日志"——一个持久化的、分布…

消息队列流处理Kafka事件驱动异步通信

2010 年,LinkedIn 的工程团队遇到了一个奇怪的问题:系统里有数十个数据管道,把数据从各个来源同步到各个目的地,形成了一张复杂的点对点网络。每增加一个新数据来源或新的数据消费者,就需要新写一批同步管道,整个系统呈现出 O(n2)O(n^2) 的连接复杂度。

他们的解决方案是发明一个"统一日志"——一个持久化的、分布式的、可重放的消息流,让所有生产者写入,所有消费者按需读取。这个系统后来被开源,名字叫 Apache Kafka,成为了现代数据基础设施的核心组件。

Kafka 于 2011 年初开源,2012 年 10 月从 Apache 孵化器毕业;2014 年三位核心作者 Jay Kreps、Neha Narkhede、Jun Rao 离开 LinkedIn 创办 Confluent,专门做 Kafka 的商业化。名字取自作家弗兰兹·卡夫卡——Kreps 解释说,这是"一个为写入而优化的系统",而他恰好喜欢卡夫卡的作品。

这个"统一日志"今天的规模,能说明它解决的问题有多大。据 LinkedIn 工程团队披露,其内部运行着 100 多个 Kafka 集群、4000 多个 broker,承载 10 万个以上的主题(Topic)、700 万个分区,每天处理超过 7 万亿条消息。

破除误解:消息队列不只是"缓冲区"

把消息队列理解为临时缓冲区——让发送方和接收方不用同时在线——这个理解只捕捉到了浅层价值。

消息队列更深层的价值在于:它把生产者和消费者在时间、空间、接口上完全解耦。生产者不知道消费者是谁,消费者不在乎生产者的实现;一个生产者的消息可以同时被十个消费者用十种方式处理;系统峰值流量可以被队列吸收,以平稳速率消费。这种解耦是构建弹性、可扩展分布式系统的基础模式之一。

两种核心范式

消息队列(Message Queue):点对点或发布-订阅模型,消息被消费后删除(或超时删除)。典型场景是任务分发——用户上传一张图片,系统把压缩任务放入队列,若干个工作进程竞争取任务处理。RabbitMQ、Amazon SQS 是代表实现。

消息流(Event Streaming):消息持久化保存,消费者可以从任意历史位置开始消费,也可以多个消费者独立消费同一条消息。Apache Kafka、Amazon Kinesis 是代表实现。流的关键区别是可重放性:数据不会因为"已被消费"而消失,可以用来重建系统状态或接入新消费者。

维度消息队列消息流
保留策略消费后删除按时间或大小保留
消费方式竞争消费(每条消息只被一个消费者处理)广播消费(每个消费者组独立读取)
顺序保证通常是全局 FIFO分区内有序,跨分区无序
典型用途任务队列、工作分发事件日志、数据集成、实时分析

Kafka 的核心设计

Kafka 的架构有几个关键决策,使其在吞吐量和可靠性之间达到了当时前所未有的平衡:

分区(Partition):每个主题(Topic)被分成多个分区,每个分区是一个有序的不可变日志。消息按写入顺序追加,偏移量(offset)单调递增。分区是并行性的单元——不同分区可以并行生产和消费,可以分布在不同节点。

消费者组(Consumer Group):一个消费者组内,每个分区只被一个消费者消费;不同消费者组独立消费所有分区。这使得"队列语义"(一条消息只处理一次)和"广播语义"(所有消费者都看到所有消息)可以同时实现,只需不同消费者组数量。

顺序写入:消息只追加到分区末尾,避免随机 I/O,充分利用磁盘顺序写的高吞吐。即使使用机械硬盘,Kafka 也能实现 GB/s 级别的写入速率。

日志压缩(Log Compaction):除了按时间或大小删除旧数据,Kafka 还提供一种基于键的保留策略——对每个消息键只保留最后一个值。这让一个主题可以充当"变更日志"(changelog):重放之后就能重建出某张表的最新快照,是事件溯源和流式状态存储的底层支撑。要删除某个键,写入一条 value 为 null 的"墓碑"(tombstone)消息即可。

不过压缩的语义有边界:它不保证任一时刻每个键只有一条记录,消费者也不保证能看到某个键的每一个历史值——只保证最终读到的是最新值。把它当成"去重的精确历史"用,会踩坑。

流处理:从批量到实时

数据处理有两种模式:

批处理(Batch Processing):积累一段时间的数据,一次性处理(如 Hadoop MapReduce)。延迟高(小时级),但吞吐量大,适合历史分析。

流处理(Stream Processing):数据一产生就立即处理,延迟毫秒到秒级。Apache Flink、Apache Spark Streaming、Kafka Streams 是主要工具。

流处理的核心挑战是时间语义

  • 事件时间(Event Time):事件实际发生的时间(嵌在消息里)
  • 处理时间(Processing Time):系统收到并处理事件的时间

网络延迟和消息乱序意味着处理时间晚于事件时间,且不稳定。基于事件时间的聚合(如"过去一分钟的点击量")需要水位线(Watermark)机制:系统估计到某个水位线时刻,该时刻之前的事件已全部到达,可以安全触发计算。Google Dataflow 论文(2015)把这套时间语义系统化,影响了之后所有流处理框架的设计。

事件驱动架构

消息队列是事件驱动架构(EDA)的物理基础。在 EDA 中,系统的状态变化以"事件"(不可变的事实记录)的形式发布,感兴趣的组件订阅并响应。

事件溯源(Event Sourcing)是一种更激进的模式:系统的当前状态不直接存储,而是通过重放所有历史事件推导出来。这让系统拥有完整的审计日志,且可以把状态"时光倒流"到任意历史时刻——代价是读取当前状态需要重放事件(通常配合快照缓解)。

Martin Kleppmann 在 2014 年 Strange Loop 大会的演讲《把数据库由内向外翻转》(Turning the Database Inside-Out)把这个思路推到极致。他指出,传统数据库把"日志"藏在内部、对外只暴露表和索引;而事件流架构正好反过来——日志成为对外的第一公民,表、索引、缓存都只是从日志派生出来的"物化视图"(materialized view),随时可以丢弃再重建。换句话说,数据库内部那条复制日志,被搬到了整个系统架构的中心。

代价与争议

顺序性的幻觉:Kafka 只保证分区内有序,跨分区的消息顺序无法保证。应用依赖全局顺序时,必须把相关消息路由到同一分区(用消息键哈希),这限制了并行度。

"至少一次"交付的隐患:为保证消息不丢,Kafka 默认提供"至少一次"(at-least-once)语义——某些情况下消息可能被重复交付。消费者必须实现幂等性(重复处理同一消息不产生副作用),这需要额外的工程成本。Kafka 0.11(2017)引入了幂等生产者和事务,才支持"恰好一次"(exactly-once)语义,但性能有所损失。

这次"恰好一次"的发布在 2017 年引发了一场不小的争论。批评者搬出两将军问题、FLP 不可能性与 CAP 定理,论证"恰好一次交付"在不可靠网络上根本无法做到。Kafka 团队(Jay Kreps 等)的回应是:这是一个语义之争——网络层的"交付次数"确实无法精确为一次,但应用层的"处理效果"可以做到一次,靠的是幂等写入加事务,把"可能重复交付、但只生效一次"落实下来。

而且这个保证既不免费、也不是端到端自动成立的。只有当下游的写入目标(sink)本身支持事务或幂等时,整条链路才真正是恰好一次;否则 Kafka 内部的恰好一次,一旦写到不支持事务的外部系统,又会退回到至少一次。

运维复杂性:长期以来 Kafka 依赖 ZooKeeper 这个独立的协调服务来管理元数据,集群管理、分区再平衡、消费者组重平衡都是棘手的运维问题。社区为此推进了 KRaft,把元数据共识内置进 Kafka 自身:2.8(2021)早期预览,3.3(2022 年 10 月)对新集群标记为生产可用,到 4.0(2025 年 3 月)彻底移除了对 ZooKeeper 的依赖。即便如此,许多团队仍转向 Confluent Cloud 等托管服务以降低运维负担。

单体 broker 的扩展之痛:Kafka 的 broker 把"计算"(处理客户端请求)和"存储"(分区数据落在本地磁盘)耦合在一起,扩容或替换节点时往往要在 broker 之间搬运大量数据,缓慢且占带宽。Apache Pulsar 给出的另一条路是把两者分离——无状态的 broker 负责服务,底层的 Apache BookKeeper(数据存为 ledger 的 bookies)负责持久化,于是存储可以独立于计算扩展,加机器不必重排数据。Kafka 社区随后也用分层存储(KIP-405,3.6 早期预览、3.9 生产可用)把冷数据卸载到对象存储,缓解了同一痛点。"存储与计算是否该解耦",至今仍是流平台架构的活跃争论。

批流之争:Lambda 与 Kappa

实时与离线如何共存,是流处理时代的一场核心架构争论。

Lambda 架构(由 Storm 之父 Nathan Marz 提出)用两条并行链路:批处理层负责准确、全量地重算历史,速度层负责低延迟的近似结果,查询时再把两者合并。它的代价很直接——同一套业务逻辑要在批和流两套系统里各写一遍、各维护一遍,两边还得保证算出来一致。

Kappa 架构(Jay Kreps 在 2014 年《质疑 Lambda 架构》一文中提出)反问:既然流处理框架已经成熟到能承担过去只有批处理能做的活,为什么还要养两套代码?Kappa 用一条可重放的日志(如 Kafka)作为唯一真相源,需要"重算"时就从日志头重新消费一遍,用同一套流式代码同时服务实时和回溯。

两者并非谁简单取代谁:当对历史结果有强准确性要求、且批流逻辑难以统一时,Lambda 仍有价值;当单一流式语义能覆盖全部需求时,Kappa 更简洁。这场争论的底层前提,正是 Kafka 那条"可重放的持久日志"——没有它,回溯重算无从谈起。

跨域连接

  • 数据库事务:日志外置是一次结构性的对调。关系数据库把预写日志藏在内部、对外只暴露表与索引;流平台把同一条日志变成公开的第一公民,于是表、索引、缓存都成了可丢弃重建的派生视图。可检验的推论是:只要保留期覆盖重建耗时,改索引、改模式、接新下游都退化成一次重放,而不再是一次迁移。
  • 信号处理:窗口聚合就是与窗函数做卷积,滚动窗、滑动窗与会话窗的区别只是窗形不同。这立刻给出取舍:窗越长,估计方差越小而延迟越大,二者不能同时优化;窗是否重叠决定同一事件被计入几次。把水位线理解成「何时认定窗内样本已到齐」,完整性与时效性同样只能在一条线上选落点。
  • 古气候与冰芯:冰芯与事件日志是同一种结构——分层累积的不可变记录,当前状态由整段记录推导而非单点读数。连压缩方式都相似:深层冰被挤压得分辨率下降,日志压缩只保留每个键的最后取值。推论一致且重要:两者都能重建现在,却重建不了被抹平的中间过程,当作精确历史来用一定出错。
  • 订单簿:交易所撮合引擎是「单分区全序日志」的极端形态:全序在这里是正确性要求而非性能选项,于是吞吐上限被单个定序器锁死。这恰好反衬出分区内有序的设计取舍——把全序放宽到分区内才换来横向扩展;反过来说,任何真正需要全局顺序的业务,都必须用并行度去付账。
  • 传染病建模与监测:公共卫生监测早就在处理事件时间与处理时间不一致的问题:发病日与报告日之间存在延迟,于是按报告日画的实时曲线末端总是向下弯,看起来像疫情正在消退。校正报告延迟与设置水位线是同一件事——都要估计还有多少已发生但未到达,并明确说明何时把窗口关上。

参考文献

  • Akidau, T. et al. The Dataflow Model. PVLDB 8(12), 2015.(Google 流处理时间语义论文,水位线机制的系统化来源)
  • Kreps, J., Narkhede, N., Rao, J. Kafka: a Distributed Messaging System for Log Processing. NetDB Workshop, 2011.(原始 Kafka 论文)
  • Kreps, J. Questioning the Lambda Architecture. O'Reilly Radar, 2014.(提出 Kappa 架构、批流之争的一手出处)
  • Kreps, J. Exactly-once Support in Apache Kafka. 2017;Confluent. Exactly-once Semantics Are Possible: Here's How Apache Kafka Does It. 2017.(恰好一次语义之争的当事方阐述)
  • Kleppmann, M. Turning the Database Inside-Out with Apache Samza. Strange Loop, 2014.(日志即第一公民、物化视图思想的来源)
  • Apache Software Foundation. KIP-405: Kafka Tiered Storage(3.6 早期预览、3.9 生产可用);KIP-833: Mark KRaft as Production Ready.(Kafka 演进的一手提案)
  • LinkedIn Engineering. How LinkedIn customizes Apache Kafka for 7 trillion messages per day. 2019.(生产规模数据一手来源)

延伸阅读

  • Kreps, J. The Log: What every software engineer should know about real-time data's unifying abstraction. engineering.linkedin.com, 2013.(Kafka 创始人写的流日志思想精华)
  • Kleppmann, M. Designing Data-Intensive Applications. Chapter 11. O'Reilly, 2017.(流处理工程实践最佳参考)
  • Narkhede, N., Shapira, G., Palino, T. Kafka: The Definitive Guide. O'Reilly, 2017.