ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

system-design-notes第19章:设计分布式消息队列完整指南

system-design-notes第19章:设计分布式消息队列完整指南 system-design-notes第19章设计分布式消息队列完整指南【免费下载链接】system-design-notesNotes of the book System Desgin Interview - An Insiders Guide项目地址: https://gitcode.com/GitHub_Trending/sy/system-design-notes本文来自system-design-notes项目《System Design Interview - An Insiders Guide》系统设计面试笔记第19章带你从零开始设计一个完整的分布式消息队列。我们将拆解 Kafka、RabbitMQ 等主流消息队列的核心机制Topic 分区、消费者组、WAL 日志存储、副本复制、Broker 故障恢复与消息投递语义是新手掌握消息队列系统设计的完整指南。为什么需要分布式消息队列消息队列是大型互联网系统的缓冲中枢引入它主要有四大好处解耦生产者和消费者不再强依赖可以独立更新、独立部署可扩展生产端与消费端可以按流量独立扩缩容️高可用某一部分宕机时其他组件仍可通过队列继续交互⚡高性能生产者发完消息即可返回无需等待消费者处理确认常见实现包括Kafka、RabbitMQ、RocketMQ、Apache Pulsar、ActiveMQ、ZeroMQ。严格来说Kafka 和 Pulsar 属于事件流平台但两者在功能上正在趋同融合。本章的目标更进一步设计一个支持两周数据保留、消息可重复消费、顺序保证的增强型消息队列——这正是传统消息队列不具备的能力。设计前先明确需求功能与非功能要求面试中先和面试官对齐需求是拿高分的关键。本章梳理出的核心约束如下类别关键要求消息格式纯文本大小在 KB 级别消费模式支持一次消费也支持被不同消费者重复消费顺序保证消息需按生产顺序被消费数据保留消息需保留两周投递语义至少支持 at-least-once理想情况下三者可配置吞吐延迟可配置日志聚合场景要高吞吐传统场景要低延迟非功能要求系统必须分布式可扩展、数据持久化落盘并在多节点间复制。消息队列核心概念速览两种消息模型点对点 vs 发布/订阅消息队列有两种经典消息传递模型。点对点模型Point-to-Point消息进入队列后被恰好一个消费者消费确认消费后即从队列删除多个消费者之间是竞争关系![分布式消息队列点对点消息模型示意](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/point-to-point-model.png?utm_sourcegitcode_repo_files)发布/订阅模型Pub/Sub消息关联到某个主题Topic订阅该主题的所有消费者都能收到完整消息副本是事件流平台的主流模型![分布式消息队列发布订阅消息模型示意](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/publish-subscribe-model.png?utm_sourcegitcode_repo_files)Topic、分区Partition与 Broker当某个 Topic 的数据量过大时扩容手段就是分区Partition即分片消息在 Topic 的各分区间均匀分布承载分区的服务器称为Broker每个分区内部是一个 FIFO 队列消息顺序在分区内得到保证消息在分区中的位置称为offset偏移量消息通过**分区键partition key**决定落到哪个分区例如用user_id作为分区键即可保证同一用户的消息有序消费者组与分区分配多个消费者可以组成**消费者组Consumer Group**共同消费一个 Topic![分布式消息队列消费者组与分区分配关系图](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/consumer-groups.png?utm_sourcegitcode_repo_files)消息是按消费者组复制的而不是按单个消费者每个组维护自己独立的 offset组内并行消费能提升吞吐但会牺牲顺序保证因此规则是一个分区同一时刻只能被组内一个消费者订阅组内消费者数量不能超过分区数分布式消息队列整体架构综合以上概念我们得到一个包含六大核心组件的整体架构![分布式消息队列整体架构图生产者、Broker、消费者与协调服务](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/high-level-architecture.png?utm_sourcegitcode_repo_files)组件职责客户端生产者Producer推送消息到 Topic消费者组Consumer Group订阅消费Broker持有多个分区是队列服务的核心节点数据存储以分区为单位存储消息状态存储保存消费者状态分区-消费者映射、消费 offset元数据存储保存配置与 Topic 属性分区数、保留期、副本分布协调服务负责服务发现哪些 Broker 存活与领导者选举设计深挖分布式消息队列的关键实现用 WAL 预写日志实现高效消息存储先分析消息数据的访问特征写多读多、无更新删除、以顺序读写为主。据此选型❌ 关系型数据库难以同时应对高写高读✅预写日志WAL, Write-Ahead Log只支持追加的纯文本文件对 HDD 极其友好分区被切分为多个段segment旧段只读只有最新段接受写入避免维护超大的单一文件![分布式消息队列WAL预写日志分段存储示例](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/wal-example.png?utm_sourcegitcode_repo_files)很多人误以为 HDD 一定慢但这取决于访问模式——顺序访问下 HDD 可达数 MB/s 的读写速度再叠加操作系统的磁盘缓存性能完全够用。消息本身采用不可变结构避免高流量场景下的额外拷贝。消息头部包含关键字段![分布式消息队列消息结构与头部字段组成](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/message-structure.png?utm_sourcegitcode_repo_files)Key决定消息归属分区如hash(key) % numPartitions与 KV 存储不同key 无需唯一甚至可以为空Value消息负载明文或压缩二进制块其余字段Topic、分区 ID、offset三者唯一定位一条消息、时间戳、大小、CRC 校验保证消息完整性生产端路由与批量写入生产者要发消息到某个分区该连接哪个 Broker一种方案是引入独立的路由层但它带来额外网络跳转、且无法批量。更优做法是把路由层内嵌到生产者内部![分布式消息队列生产者内嵌路由层与批量缓冲设计](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/routing-layer-producer.png?utm_sourcegitcode_repo_files)生产者本地缓存复制计划直接连接分区领导者减少一跳内存缓冲区支持**批量Batching**发送把多条消息攒成大批次一次性写入 WAL 的顺序写摊薄网络与磁盘开销大幅提升吞吐批量大小是经典的吞吐-延迟权衡批量越大吞吐越高但延迟越高反之延迟低但吞吐低。若按低延迟场景部署调小批量即可若按高吞吐调优则需要更多分区来弥补单分区顺序写的速度瓶颈。消费者拉取模型与重平衡消费端采用指定 offset 拉取消息的方式。选型时最关键的是Push 还是 PullPush 模型延迟低但消费慢时消费者会被冲垮且 Broker 难以适配处理能力强弱不一的消费者Pull 模型消费速率由消费者自己掌控可扩容追赶适合批处理代价是空轮询带来的额外请求可用**长轮询long polling**缓解因此绝大多数消息队列包括本设计都选择 Pull 模型![分布式消息队列消费者拉取消息与提交offset流程](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/consumer-flow.png?utm_sourcegitcode_repo_files)消费者加入分区的完整流程新消费者订阅 Topic 并申请加入某消费者组通过对组名做哈希定位负责该组的 Broker即组协调器注意它和 ZooKeeper 协调服务不是一回事协调器确认入组并分配分区支持轮询、范围等多种分配策略消费者从状态存储中读取上次的 offset开始拉取最新消息处理完成后向 Broker 提交 offset——处理与提交 offset 的先后顺序直接决定了投递语义当有消费者加入/离开、或分区数量变化时会触发消费者重平衡Rebalancing。同一组的所有消费者都连接到同一个协调 Broker![分布式消息队列消费者组重平衡协调机制](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/consumer-rebalancing.png?utm_sourcegitcode_repo_files)成员列表变化后协调器为该组选举新的组内领导者领导者计算新的分区分配方案并上报协调器广播给组内所有消费者当协调器长时间收不到某消费者的心跳消费者宕机同样会触发重平衡保证故障容忍用 ZooKeeper 管理元数据与消费者状态消费者状态分区-消费者映射、各分区最后消费的 offset具有高频低量、随机读写、强一致的访问特征非常适合快速的 KV 存储元数据分区数、保留期、副本分布变更不频繁但要求强一致。两者都交由ZooKeeper承担![分布式消息队列中ZooKeeper存储元数据与协调服务](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/zookeeper.png?utm_sourcegitcode_repo_files)这样改造后Broker 只专注存储消息数据元数据与状态全部交给 ZooKeeper它同时协助 Broker 副本的领导者选举。副本机制与 ISR高可用的基石硬件故障不可避免必须靠副本复制实现高可用。每个分区复制在多台 Broker 上其中只有一台是领导者![分布式消息队列分区副本跨Broker复制拓扑](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/replication-example.png?utm_sourcegitcode_repo_files)生产者只向领导者副本写入跟随者副本主动向领导者拉取消息足够多的副本同步后领导者才向生产者返回确认各分区的副本分布方案由领导者制定并保存进 ZooKeeper为了判断副本是否跟得上引入ISRIn-Sync Replicas同步副本集合滞后超过阈值如replica.lag.max.messages的副本会被移出 ISR![分布式消息队列ISR同步副本与committed offset示例](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/in-sync-replicas-example.png?utm_sourcegitcode_repo_files)确认策略ACK可配置直接反映性能与持久性的权衡配置行为特点ACKall等待 ISR 内所有副本同步最慢但持久性最高ACK1领导者收到即确认速度快持久性较低ACK0发出不等待任何确认最快可能丢消息![分布式消息队列ACKall确认机制下副本同步流程](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/ack-all.png?utm_sourcegitcode_repo_files)消费侧可以让所有消费者都连到分区领导者读取分区内消息同一时刻只发给组内一个消费者连接数天然可控超热的 Topic 则通过增加分区与消费者来横向扩展。跨数据中心场景下也可以让消费者就近从 ISR 副本读取。Broker 故障恢复与副本再均衡当某台 Broker 宕机时系统依靠副本自动恢复无需人工干预![分布式消息队列Broker故障恢复与重新选主流程图](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/broker-failure-recovery.png?utm_sourcegitcode_repo_files)Broker-3 故障后其分区的其他副本依然保有数据不会丢失重新选举领导者协调器把故障 Broker 上的分区重新分配给存活副本新副本先作为跟随者追赶进度追平后进入 ISR工程上还要注意副本应跨不同 Broker甚至跨机房打散ISR 最小数量决定了延迟与安全性的平衡点可按业务微调。扩容时新 Broker 加入后允许临时超配副本数待其追平数据后再移除多余副本缩容分区时退役分区不会被立即删除——生产者只向活跃分区写入消费者继续读完存量消息等保留期过期后再截断释放空间并重平衡消费者。三种消息投递语义如何选择投递语义是面试高频考点三者权衡一目了然语义生产者行为消费者行为特点至多一次At-most-once异步发送失败不重试拉取后立即提交 offset可能丢消息绝不重复至少一次At-least-onceack1/all失败持续重试处理完成后才提交 offset可能重复不丢消息适合可去重场景精确一次Exactly-once——对用户最友好但实现成本极高核心陷阱在于消费者已处理消息但崩溃在提交 offset 之前重平衡后新消费者会重放该消息——这就是 at-least-once 产生重复的根源业务侧需做好幂等或去重。进阶特性消息过滤与延迟消息消息过滤若某类消费者只想消费分区中特定类型的消息为每种需求单独建 Topic 代价太高重复存储、生产者与消费者强耦合。优雅方案是给消息打标签tag消费者声明订阅哪些标签由 Broker 侧完成过滤避免把无关流量打到消费者端![分布式消息队列消息标签过滤机制示意](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/message-filtering.png?utm_sourcegitcode_repo_files)延迟与定时消息典型场景发起支付后 30 分钟再触发消费者检查支付是否成功。实现方式是先把消息写入 Broker 的临时存储到期后再移入正式分区![分布式消息队列延迟消息定时投递实现方案](https://raw.gitcode.com/GitHub_Trending/sy/system-design-notes/raw/9d8388721e7231442763ad37398b8d82224aa68f/19. Distributed Message Queue/images/delayed-message-implementation.png?utm_sourcegitcode_repo_files)定时调度可以用专用延迟队列也可以采用**分层时间轮Hierarchical Time Wheel**等高效算法。本章总结与延伸阅读回顾本章设计分布式消息队列的完整脉络✅需求先行保留期、顺序、重复消费等增强需求决定了架构走向✅存储选型WAL 顺序追加 分段 不可变消息结构吃透 HDD 顺序读写优势✅批量写入生产端内嵌路由与内存缓冲吞吐与延迟动态权衡✅消费者体系Pull 模型 消费者组 重平衡兼顾吞吐、顺序与故障容忍✅高可用副本复制 ISR 可配置 ACK 策略Broker 故障自动选主恢复✅可扩展生产、消费、Broker、分区四个维度均可独立水平扩展更多细节如通信协议选型、重试消费、历史数据归档到 HDFS/对象存储等可以查阅本章完整原文19. Distributed Message Queue/README.md项目整体章节索引见 Readme.md。掌握这套设计范式后你再去理解 Kafka 等开源消息队列的源码与文档就会事半功倍 【免费下载链接】system-design-notesNotes of the book System Desgin Interview - An Insiders Guide项目地址: https://gitcode.com/GitHub_Trending/sy/system-design-notes创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进