第 19 章:分布式消息队列
引言
本章设计一个分布式消息队列。
消息队列的好处:
- 解耦:消除组件之间的紧密依赖,让它们独立更新。
- 提高可扩展性:可根据流量分别扩展生产者和消费者。
- 提高可用性:系统某一部分故障时,其他部分仍可与队列交互。
- 改善性能:生产者无需等待消费者确认,就能继续生产消息。
常见实现包括 Kafka、RabbitMQ、RocketMQ、Apache Pulsar、ActiveMQ 和 ZeroMQ。
严格来说,Kafka 和 Pulsar 是事件流平台,并非消息队列。不过,两类系统的功能逐渐融合,界限变得模糊。
本章设计的消息队列支持长期保留数据、重复消费消息等较高级的功能。
第 1 步:理解问题并确定设计范围
消息队列的基本功能是让生产者生产消息、消费者消费消息。但性能、消息投递和数据保留等方面仍有不同选择。
候选人与面试官可能会这样讨论:
- 候选人:消息是什么格式、平均多大?仅包含文本吗?
- 面试官:只有文本,通常为几 KB。
- 候选人:消息能重复消费吗?
- 面试官:能,不同消费者可以重复消费;这是传统消息队列不支持的额外需求。
- 候选人:消费顺序要与生产顺序一致吗?
- 面试官:是,需要保证顺序;这也是额外需求。
- 候选人:数据要保留多久?
- 面试官:两周;同样是额外需求。
- 候选人:需要支持多少生产者和消费者?
- 面试官:越多越好。
- 候选人:投递语义要支持至多一次、至少一次还是恰好一次?
- 面试官:至少要支持“至少一次”;最好全部支持,并且可配置。
- 候选人:吞吐量和端到端延迟有什么目标?
- 面试官:日志聚合等场景需要高吞吐量,传统场景则可能需要低延迟。
功能需求
- 生产者向消息队列发送消息。
- 消费者从队列消费消息。
- 消息可以消费一次或重复消费。
- 可以清理历史数据。
- 消息大小为 KB 级。
- 需要保持消息顺序。
- 投递语义可配置:至多一次、至少一次、恰好一次。
非功能需求
- 高吞吐量或低延迟:根据使用场景配置。
- 可扩展:系统应采用分布式架构,能应对消息量骤增。
- 持久性与耐久性:数据应写入磁盘,并在节点之间复制。
传统消息队列通常不保留数据,也不保证顺序,这会大幅简化设计;下文将讨论这些差异。
第 2 步:提出概要设计并达成共识
消息队列的关键组件:
<img src="/images/chapter-19/message-queue-components.png" alt="消息队列组件" width="500" />
- 生产者向队列发送消息。
- 消费者订阅队列并消费订阅的消息。
- 消息队列位于两者之间,使生产者与消费者解耦并独立扩展。
- 生产者和消费者都是客户端,消息队列是服务器。
消息传递模型
第一种是点对点模型,传统消息队列中很常见:
<img src="/images/chapter-19/point-to-point-model.png" alt="点对点模型" width="500" />
- 消息发送到队列后,由恰好一个消费者消费。
- 可以有多个消费者,但每条消息只消费一次。
- 消息确认已消费后,从队列移除。
- 点对点模型通常不保留数据,但本章的设计会保留。
另一种是发布/订阅模型,事件流平台更常使用:
<img src="/images/chapter-19/publish-subscribe-model.png" alt="发布订阅模型" width="500" />
- 消息归属于某个主题。
- 消费者订阅主题,接收发往该主题的全部消息。
主题、分区和 Broker
如果某个主题的数据量太大,可以将其拆成多个分区,也就是分片:
<img src="/images/chapter-19/partitions.png" alt="分区" width="500" />
- 发往主题的消息均匀分布到多个分区。
- 承载分区的服务器称为 Broker。
- 每个分区像 FIFO 队列一样处理消息;顺序在分区内得到保证。
- 消息在分区中的位置称为偏移量(offset)。
- 每条消息会进入特定分区;分区键决定目标分区。
- 例如,使用
user_id作为分区键,可保证同一用户的消息顺序。
- 例如,使用
- 每个消费者订阅一个或多个分区。若多个消费者协同处理同一批消息,它们组成消费者组。
消费者组
消费者组是一组协同消费某个主题消息的消费者:
<img src="/images/chapter-19/consumer-groups.png" alt="消费者组" width="500" />
- 消息按消费者组复制,而非按单个消费者复制。
- 每个消费者组维护自己的偏移量。
- 组内并行读取可以提高吞吐量,但可能破坏顺序保证。
- 限制同一分区在一个组内只由一位消费者订阅,可以缓解这个问题。
- 因此,一个组内的消费者数量不能超过分区数量。
概要架构
<img src="/images/chapter-19/high-level-architecture.png" alt="概要架构" width="500" />
- 客户端:包括生产者与消费者。生产者向指定主题发送消息,消费者组订阅主题消息。
- Broker:承载多个分区,每个分区存放主题的部分消息。
- 数据存储:存储各分区中的消息。
- 状态存储:保存消费者状态。
- 元数据存储:保存配置和主题属性。
- 协调服务:负责服务发现(哪些 Broker 存活)和选主(哪个 Broker 负责分配分区)。
第 3 步:深入设计
为实现高吞吐量和较长的数据保留期,设计中做出几项重要选择:
- 选择能利用现代 HDD 特性和操作系统磁盘缓存策略的磁盘数据结构。
- 消息数据结构不可变,以避免高流量场景中的额外复制。
- 围绕批量处理设计写入流程,因为大量小 I/O 会损害吞吐量。
数据存储
选择消息存储方式前,先看消息的特性:
- 读写量都很大。
- 不更新、不删除。传统消息队列因不保留消息而有“删除”操作。
- 主要采用顺序读写。
可选方案:
- 数据库:通常难以同时高效承担大量读写,并不理想。
- 预写日志(WAL):只支持追加的普通文件,很适合 HDD。
把分区拆为多个段,避免维护过大的单个文件。
旧段只读,仅向最新段写入。

WAL 在传统 HDD 上非常高效。
认为 HDD 访问一定慢是个误解,实际速度高度依赖访问模式。顺序访问时,HDD 可达到每秒数 MB 的读写速度,足以满足本系统需求。操作系统还会积极把磁盘数据缓存到内存。
消息数据结构
生产者、队列和消费者应使用一致的消息结构,避免额外复制,从而提高处理效率。
消息结构示例:
<img src="/images/chapter-19/message-structure.png" alt="消息结构" width="500" />
消息键指定其所属分区,例如 hash(key) % numPartitions。生产者也可以覆盖默认键,控制消息分配到哪个分区。
消息值是载荷,可以是纯文本或压缩后的二进制块。
注意: 与传统键值存储不同,消息键无需唯一,可以重复,甚至可以缺失。
其他消息字段:
- 主题(Topic):消息所属的主题。
- 分区(Partition):消息所属分区的 ID。
- 偏移量(Offset):消息在分区中的位置,可用
topic、partition、offset定位消息。 - 时间戳(Timestamp):消息的存储时间。
- 大小(Size):消息大小。
- CRC:用于验证消息完整性的校验和。
可通过增加字段支持过滤等附加功能。
批量处理
批量处理对系统性能至关重要,应在生产者、消费者和消息队列中使用。
原因是:
- 操作系统可以合并消息,分摊昂贵的网络往返成本。
- 消息按组顺序写入 WAL,从而获得大量顺序写入及磁盘缓存收益。
延迟与吞吐量之间需要权衡:
- 批次越大,吞吐量越高,但延迟也越高。
- 批次越小,吞吐量越低,但延迟也越低。
若系统作为传统消息队列,需要较低延迟,就可以调小批次。
若系统针对吞吐量优化,则可能需要为每个主题增加分区,以弥补顺序磁盘写入速度的限制。
生产者流程
生产者向某个分区发送消息时,应连接哪台 Broker?
一种选择是增加路由层,由它将消息送到正确的 Broker。若启用了复制,目标 Broker 就是主副本所在的 Broker:
<img src="/images/chapter-19/routing-layer.png" alt="路由层" width="500" />
- 路由层从元数据存储读取副本分布计划,并在本地缓存。
- 生产者向路由层发送消息。
- 消息转发给该分区的主副本,即 Broker 1。
- 从副本向主副本拉取新消息。收到足够的确认后,主副本提交数据并回复生产者。
使用副本是为了容错。
这种方式可行,但有缺点:
- 额外的组件会增加网络跳数。
- 不便于对消息进行批量处理。
为解决这些问题,可以把路由层嵌入生产者:
<img src="/images/chapter-19/routing-layer-producer.png" alt="嵌入生产者的路由层" width="500" />
- 网络跳数减少,延迟降低。
- 生产者可控制消息进入哪个分区。
- 缓冲区让生产者在内存中组成批次,用一次请求发送更大的消息批,提高吞吐量。
批次大小是吞吐量与延迟之间的经典权衡。
<img src="/images/chapter-19/batch-size-throughput-vs-latency.png" alt="批次大小、吞吐量与延迟" width="500" />
- 批次越大,提交前等待越久。
- 批次越小,请求发送越早,延迟越低,但吞吐量也越低。
消费者流程
消费者指定某分区的偏移量,从该位置开始获取一批消息:
<img src="/images/chapter-19/consumer-example.png" alt="消费者示例" width="500" />
设计消费者时,一个重要选择是推送还是拉取:
- 推送模型:Broker 收到消息后主动推送给消费者,延迟较低。
- 但如果消费速度落后于生产速度,消费者可能被压垮。
- Broker 控制消费速度,难以适配处理能力不同的消费者。
- 拉取模型:消费者自行控制消费速度。
- 消费慢时不会被压垮,并可扩容追赶积压。
- 更适合批处理;推送模型下 Broker 不知道消费者一次能处理多少消息。
- 拉取模型中,消费者可以主动获取大量消息。
- 缺点是在没有新消息时延迟较高、产生额外网络请求;长轮询可缓解后者。
因此,多数消息队列(包括本设计)选择拉取模型。
<img src="/images/chapter-19/consumer-flow.png" alt="消费者流程" width="500" />
- 新消费者订阅主题 A,并加入组 1。
- 对组名做哈希,找到对应的 Broker 节点;同组消费者都连接该节点。
- 这里的消费者组协调器不同于 ZooKeeper 协调服务。
- 协调器确认消费者入组,并将分区 2 分配给它。
- 分区分配策略包括轮询、范围分配等。
- 消费者从上次偏移量开始获取最新消息。状态存储维护消费者偏移量。
- 消费者处理消息,并向 Broker 提交偏移量。两项操作的顺序决定消息投递语义。
消费者再均衡
再均衡决定每个分区由哪些消费者负责。
消费者加入、离开,或分区增加、移除时,都会触发再均衡。
作为协调器的 Broker 负责组织再均衡流程。
<img src="/images/chapter-19/consumer-rebalancing.png" alt="消费者再均衡" width="500" />
- 同组消费者连接同一个协调器;通过组名哈希找到协调器。
- 消费者列表变化后,协调器选择新的组负责人。
- 负责人计算新的分区分配计划并上报协调器,协调器再广播给其他消费者。
协调器长时间收不到组内消费者的心跳,也会触发再均衡:
<img src="/images/chapter-19/consumer-rebalance-example.png" alt="消费者再均衡示例" width="500" />
消费者入组时的流程:
<img src="/images/chapter-19/consumer-join-group-usecase.png" alt="消费者加入组" width="500" />
- 起初组内只有消费者 A,负责所有分区。
- 消费者 B 请求入组。
- 协调器通过对心跳的响应,通知所有组员开始再均衡。
- 所有消费者重新入组后,协调器选出负责人并公布结果。
- 负责人生成分区分配计划并发给协调器;其他消费者等待计划。
- 消费者开始从新分配的分区消费。
消费者离组时的流程:
<img src="/images/chapter-19/consumer-leaves-group-usecase.png" alt="消费者离开组" width="500" />
- 消费者 A、B 位于同一组。
- 消费者 B 请求离组。
- 协调器收到 A 的心跳时,通知它需要再均衡。
- 后续步骤与入组流程相同。
消费者长时间不发送心跳时,流程也类似:
<img src="/images/chapter-19/consumer-no-heartbeat-usecase.png" alt="消费者没有心跳" width="500" />
状态存储
状态存储保存分区与消费者的映射,以及每个分区最后消费的偏移量。
<img src="/images/chapter-19/state-storage.png" alt="状态存储" width="500" />
组 1 的偏移量为 6,表示此前的消息均已消费。消费者崩溃后,新的消费者将从该位置继续。
消费者状态的数据访问模式:
- 读写频繁,但数据量小。
- 更新频繁,删除很少。
- 随机读写。
- 对一致性要求较高。
因此,ZooKeeper 之类的快速键值存储很合适。
元数据存储
元数据存储保存配置和主题属性,包括分区数量、保留期限及副本分布。
元数据量小且变化不频繁,但需要强一致性,因此 ZooKeeper 是合适选择。
ZooKeeper
ZooKeeper 是构建分布式消息队列的重要组件。
它是层次化键值存储,常用于分布式配置、同步服务和命名注册(即服务发现)。
<img src="/images/chapter-19/zookeeper.png" alt="ZooKeeper" width="500" />
采用这一设计后,Broker 只需维护消息数据;元数据和状态由 ZooKeeper 保存。
ZooKeeper 还能协助 Broker 副本选主。
复制
分布式系统中硬件故障不可避免。通过复制可实现高可用性。
<img src="/images/chapter-19/replication-example.png" alt="复制示例" width="500" />
- 每个分区跨多个 Broker 复制,但只有一个主副本。
- 生产者向主副本发送消息。
- 从副本从主副本拉取消息。
- 足够多的副本完成同步后,主副本向生产者返回确认。
- 每个分区的副本分布称为副本分布计划。
- 给定分区的主副本制定分布计划,并存入 ZooKeeper。
同步副本
需要解决的问题是让一个分区的主副本与从副本保持同步。
同步副本(ISR)指与主副本保持同步的分区副本。
replica.lag.max.messages 定义副本最多可以落后主副本多少条消息,仍被视为同步。
<img src="/images/chapter-19/in-sync-replicas-example.png" alt="同步副本示例" width="500" />
- 已提交偏移量为 13。
- 两条新消息已写入主副本,但尚未提交。
- ISR 中所有副本都同步某条消息后,该消息才会提交。
- 副本 2、3 已追上主副本,因此属于 ISR。
- 副本 4 已落后,因此暂时从 ISR 移除。
ISR 体现了性能与耐久性之间的权衡:
- 若要防止生产者丢失消息,发送确认前应确保所有副本同步。
- 但缓慢的副本可能使整个分区不可用。
确认策略可配置。
ACK=all 表示 ISR 内所有副本都必须同步消息;发送较慢,但耐久性最高。
<img src="/images/chapter-19/ack-all.png" alt="ACK=all" width="500" />
ACK=1 表示主副本收到消息后,生产者就获得确认;发送较快,但耐久性较低。
<img src="/images/chapter-19/ack-1.png" alt="ACK=1" width="500" />
ACK=0 表示生产者发送消息后不等待主副本确认;发送最快,但耐久性最低。
<img src="/images/chapter-19/ack-0.png" alt="ACK=0" width="500" />
在消费者侧,可以让所有消费者连接分区主副本并从中读取:
- 设计和运维最简单。
- 一个分区的消息在每组中只发送给一个消费者,因此主副本连接数有限。
- 只要主题不是极热,主副本连接数通常不会太高。
- 可以增加分区和消费者数量来扩展热门主题。
- 某些情况下可允许消费者从 ISR 读取,例如消费者位于另一数据中心时。
主副本跟踪自己与各副本之间的滞后,维护 ISR 列表。
可扩展性
下面评估系统各部分的扩展方式。
生产者
生产者比消费者简单得多。增加或移除生产者实例即可扩缩容。
消费者
消费者组彼此隔离,可以按需增加或移除。
消费者加入或离组时,再均衡机制可以妥善处理。
消费者组及其再均衡机制为系统提供可扩展性和容错能力。
Broker
Broker 故障如何处理?
<img src="/images/chapter-19/broker-failure-recovery.png" alt="Broker 故障恢复" width="500" />
- 一台 Broker 故障后,其他副本仍能避免分区数据丢失。
- 选出新主副本,Broker 协调器把故障 Broker 的分区重新分配给现有副本。
- 现有副本接管新分区并作为从副本运行,直到追上主副本并加入 ISR。
增强 Broker 容错能力还要考虑:
- ISR 最小数量在延迟和安全性之间权衡,可按需求调整。
- 若一个分区的全部副本都在同一节点,复制就失去意义;应分布到不同 Broker。
- 如果分区的全部副本都崩溃,数据将永久丢失。跨数据中心分布副本可缓解,但会显著增加延迟。可考虑使用数据镜像。
增加新 Broker 时,如何重新分配副本?
<img src="/images/chapter-19/broker-replica-redistribution.png" alt="Broker 副本重新分配" width="500" />
- 新 Broker 追上进度前,可暂时允许副本数量超过配置值。
- 追上后再移除不再需要的分区副本。
分区
新增分区时,通知生产者并触发消费者再均衡。
存储方面,可以让新分区只接收新消息,不复制全部历史消息:
<img src="/images/chapter-19/partition-exmaple.png" alt="分区示例" width="500" />
减少分区数量更复杂:
<img src="/images/chapter-19/partition-decrease.png" alt="减少分区" width="500" />
- 分区停用后,新消息只进入其余分区。
- 停用的分区不会立刻删除,因为仍可能有消息需要消费。
- 预设保留期限过后,清理数据并释放存储空间。
- 过渡期内,生产者只向活跃分区发送消息,消费者则读取所有分区。
- 保留期限届满后,对消费者再均衡。
数据投递语义
下面讨论不同的投递语义。
至多一次
消息最多投递一次,也可能完全没有投递。
<img src="/images/chapter-19/at-most-once.png" alt="至多一次" width="500" />
- 生产者异步向主题发送消息;投递失败也不重试。
- 消费者获取消息后立即提交偏移量。如果处理前崩溃,消息就不会被处理。
至少一次
消息可能投递多次,但不应有消息未被处理。
<img src="/images/chapter-19/at-least-once.png" alt="至少一次" width="500" />
- 生产者以
ack=1或ack=all发送消息,发生问题则持续重试。 - 消费者获取消息,处理完成后才提交偏移量。
- 若消费者处理完消息、提交偏移量前崩溃,消息可能重复投递。
- 因此,它适用于允许重复数据或可执行去重的场景。
恰好一次
这是对用户最友好的保证,但在系统中实现的成本极高:
<img src="/images/chapter-19/exactly-once.png" alt="恰好一次" width="500" />
高级功能
下面讨论面试中可能涉及的高级功能。
消息过滤
有些消费者只希望接收分区内某种类型的消息。
可以为不同消息子集分别建主题,但如果使用场景很多,成本会很高:
- 同一消息存放在不同主题,浪费资源。
- 每增加一种消费者需求,生产者都要改动,造成紧密耦合。
可以用消息过滤解决:
最简单的办法是在消费者侧过滤,但会带来不必要的传输。
也可以给消息附加标签,由消费者指定订阅哪些标签。
还可以按消息载荷过滤,但对于加密或序列化的消息,这可能困难且不安全。
对复杂数学表达式,Broker 可实现语法解析器或脚本执行器,但会让消息队列过于沉重。

延迟消息与定时消息
有些场景需要延迟或定时投递。例如,安排 30 分钟后检查一次付款,让消费者确认付款是否成功。
可以先把消息放入 Broker 的临时存储,再在指定时间将其移入分区:
<img src="/images/chapter-19/delayed-message-implementation.png" alt="延迟消息实现" width="500" />
- 临时存储可以是一个或多个特殊消息主题。
- 计时功能可用专用延迟队列或分层时间轮实现。
第 4 步:总结
还可讨论以下内容:
- 通信协议:应支持全部使用场景和高数据量,并能验证消息完整性。常见协议包括 AMQP 和 Kafka 协议。
- 消费重试:若消息暂时无法处理,可将其送入专用重试主题,稍后再次尝试。
- 历史数据归档:旧消息可以备份到 HDFS 或对象存储(例如 S3)等大容量存储中。