从零设计一个消息中间件:高性能、高可用、数据不丢失

发布时间:2026/10/11 13:41:41
从零设计一个消息中间件:高性能、高可用、数据不丢失 技术分享 · 消息中间件架构设计核心命题如果让我从零设计一个消息中间件应该如何一步步推导出它的存储结构、消费模型、分布式架构与可靠性机制一、设计目标与核心问题消息中间件的本质很简单生产者生产消息消费者消费消息中间件负责消息的存储、路由与投递。但一旦加上生产级约束——日均 TB 级写入、机器随时会宕机、一条消息都不能丢——设计立刻变成一个系统工程。我把要解决的问题域拆成五类每一类对应本文的一章问题域核心问题对应章节存储消息放内存还是磁盘文件怎么组织第二章消费消息怎么投递给消费者怎么分配第三章分布式TB 级数据怎么存机器挂了怎么办第四章可靠性怎么保证一条消息都不丢第五章性能怎么把吞吐做到极致第六章三个设计关键词贯穿全文高性能吞吐、高可用容错、数据不丢失可靠性。1.1 总体架构蓝图架构分层自上而下接入层生产者/消费者通过协议接入路由到分区 Leader 所在 Broker服务层Broker 集群负责消息接收、存储、副本同步与消费投递Controller 负责元数据与故障转移存储层每个分区对应一组 Segment 文件内存页缓存做缓冲、磁盘做持久化协调层消费位点管理与消费者组的 Rebalance 协调。二、存储层设计内存缓冲 磁盘顺序写2.1 介质选型内存做缓冲磁盘做持久化第一个决策是消息放哪里。三种方案对比方案优点缺点结论纯内存如 Redis读写极快宕机即丢数据容量受限于内存价格不可接受消息需要持久化纯磁盘同步写可靠每条消息一次 fsync性能极差不可接受内存缓冲 磁盘持久化混合写内存立即返回吞吐高磁盘保证持久性需要处理刷盘时机与宕机窗口采用我的选择是消息先写入内存作为缓冲再定期批量刷入磁盘文件。这里的内存不是自己在 JVM 堆里维护一个缓冲区而是直接利用操作系统页缓存Page Cache顺序追加写入文件时数据先落在页缓存中由 OS 异步刷盘。好处有三避免 JVM GC 压力堆外由 OS 管理Broker 重启后页缓存仍然有效只要没被置换热数据读取接近内存速度消费者读的往往是刚写入的消息直接命中页缓存读磁盘的次数很少。2.2 写入路径顺序追加 批量刷盘这里的关键洞察是磁盘顺序写的速度远超随机写。机械盘顺序写可达 600MB/s 以上随机写只有约 100KB/s相差数千倍SSD 上差距缩小但顺序写依然占优。所以我规定存储模型只有一种写模式——在文件尾部顺序追加Append-Only绝不原地修改。这是整个设计性能的基石。2.3 文件分段Segment 机制磁盘文件不能只有一个必须按规则拆分成多个Segment段文件原因有二单文件无限增长后无法管理删除过期数据、文件句柄、索引加载按段滚动后过期数据的清理退化为删除整个旧段文件代价是 O(1) 的不需要复杂的空间回收。滚动条件按大小或时间二选一先到者为准例如段文件达到 1GB 或当前段打开超过 7 天就滚动新段。段文件以起始 offset 命名天然有序topic-order-0/ # Topic order 的 Partition 0 ├── 00000000000000000000.log # 段 1存 offset [0, 5367850] 的消息 ├── 00000000000000000000.index # 段 1 的稀疏索引 ├── 00000000000005367851.log # 段 2从 offset 5367851 开始 ├── 00000000000005367851.index ├── 00000000000010293455.log # 段 3当前活跃段 ├── 00000000000010293455.index └── leader-epoch-checkpoint # Leader 纪元快照用于副本一致性保留策略同样按时间或大小配置如保留 7 天 / 每分区最多 100GB后台线程周期性检查把整体过期的旧段直接删掉。用时间命名的维度还带来一个重要能力消息回溯——消费位点指到 3 天前的 offset就能重放 3 天前的数据这是内存队列做不到的。2.4 消息元数据offset 与消息格式每条消息必须有元数据最核心的是offset分区内单调递增的 64 位偏移量它同时承担三个职责——消息的唯一寻址 ID、消费进度位点、按范围拉取的依据。完整消息格式设计如下字段长度说明offset8B分区内唯一、单调递增消息的身份标识size4B消息总长度CRC324B校验和检测磁盘/网络传输中的消息损坏magic1B格式版本号支持消息格式平滑演进attributes1B标志位压缩算法、消息类型等timestamp8B写入时间戳支撑按时间保留与回溯key变长带 4B 长度前缀分区路由键同一 key 恒定路由到同一分区value变长带 4B 长度前缀消息体可整体压缩2.5 稀疏索引消息定位流程为每条消息建稠密索引太贵。采用稀疏索引只在段内每写入 4KB约才添加一条索引项记录offset → 物理文件位置的映射。索引文件足够小可以整个 mmap 进内存。消息定位流程二分查找 O(log n) 定位 顺序扫描少量消息整个读路径没有一次随机 IO数据在页缓存中时更是纯内存操作。2.6 零拷贝消费端读取的加速器传统读取链路要经历 4 次拷贝、4 次上下文切换磁盘 → 内核缓冲 → 用户空间 → Socket 缓冲 → 网卡。而消费是典型的读取后原样发送场景根本不需要看到内容所以引入零拷贝sendfile链路拷贝次数上下文切换用户态参与传统 read write4 次4 次是数据白白路过应用sendfile 零拷贝2 次DMA2 次否内核直接页缓存 → 网卡配合 2.1 节的页缓存策略消费热点数据时既不读磁盘、也不拷贝到用户态这是 Kafka 单机支撑极高消费吞吐的核心手段之一。三、消费模型设计消息落盘之后第二个大问题怎么投递给消费者3.1 Push 还是 Pull两大经典实现给出了不同答案对比参考维度KafkaPull 拉模式RabbitMQPush 推模式发起方消费者主动 FetchBroker 主动推送流量控制天然自适应消费者按自己的处理能力拉依赖 prefetchQoS限流推太快会压垮慢消费者消费速率消费者可控空闲时可有长轮询兜底Broker 视角统一调度状态管理位点由消费者提交Broker 无逐条消息状态Broker 维护 unacked 队列宕机可重投适合场景高吞吐、流式/回溯消费低延迟、精细的逐条确认我的选型以 Pull 为主模型吞吐、流控、回溯都占优同时吸收 RabbitMQ 的两个优点——长轮询没有新消息时 Broker 挂起请求几十毫秒再返回避免空轮询延迟逼近推模式和服务端逐条 ack 状态的可靠性思想融合进第五章的消费确认机制。3.2 消费者组组内竞争组间广播消费以消费者组Consumer Group为单位组织这是消费模型的核心抽象两条铁律组内竞争一个分区在同一时刻只分配给组内一个消费者保证组内消息只被处理一次这是水平扩展消费能力的方式组间广播每个组独立消费全量消息、各自维护位点互不影响订单服务和对账服务各建一个组即可各拿到全量消息。推论消费者数量超过分区数时多出的消费者空闲——分区数就是单组消费并行度的上限这是容量规划的关键参数。3.3 分区分配与 Rebalance消费者加入、退出、宕机或分区数变化时需要重新分配分区即Rebalance三种经典分配策略的取舍策略思路缺点Range按区间按.topic 依次切区间分给消费者多 Topic 时前面消费者总是多分倾斜RoundRobin轮询所有分区排序后逐个轮转分配Rebalance 后分区大面积换主人状态缓存失效Sticky粘性尽量保留上次分配只动必须动的实现复杂但换主最少演进版 CooperativeSticky 还支持增量 Rebalance避免 stop-the-worldRebalance 是消费稳定性的最大敌人期间整组停止消费所以工程上要合理设置心跳超时、避免频繁上下线、优先增量式协调。3.4 消费位点offset管理消费者自己维护消费进度定期把 offset提交到内部主题如__consumer_offsets持久化。重启后从上次提交的位点继续。提交时机直接决定语义与可靠性详见 5.4自动定时提交省事但处理完之前就提交了会丢消息处理成功后手动提交可靠我的默认选择。3.5 投递语义语义含义实现方式风险At most once至多一次先提交位点再处理宕机可能丢消息At least once至少一次先处理成功再提交位点宕机恢复后重复消费Exactly once精确一次至少一次 幂等生产端幂等/事务 消费端幂等成本高通常只在关键链路开启我的选型默认 At least once在业务侧做幂等达到效果上的 Exactly once唯一键约束、去重表、状态机判断。这是吞吐、复杂度与可靠性之间最均衡的点。四、分布式架构设计4.1 数据分片Partition单机存不下日均几百 TB 的数据也撑不住对应的读写压力。答案是把 Topic 水平切分成多个Partition分区均匀分布到集群各台 Broker 上每个分区是一个独立的、不可修改的追加日志第二章的存储模型以分区为单位实例化topic-order-0/、topic-order-1/……分区数决定了并行度生产端可并行写入消费端并行拉取带 key 的消息按哈希路由到固定分区分区内保持有序全局有序需要单分区吞吐换顺序按业务权衡。4.2 多副本冗余Leader/Follower 与 ISR分布式之后任何一台机器都可能宕机。所以每个分区配置多个副本Replica分散在不同 Broker 上机架感知甚至可以跨机架/机房分布。副本间角色分工要点每个分区只有一个 Leader 副本承担全部读写其余 Follower 唯一的职责就是向 Leader 发 Fetch 请求拉取同步与消费者拉取复用同一套机制设计上极简ISRIn-Sync Replicas同步副本集合跟得上 Leader 进度的副本集合落后超过阈值即被踢出。acksall 的all指的是 ISR 全体而不是所有副本——避免一个慢副本拖垮写入Follower 落后会被移出 ISR追上后再加回来。数据是否算安全由 ISR 决定这个动态集合就是一致性与可用性的调节阀。4.3 高可用Leader 选举与故障转移Leader 所在 Broker 宕机时由Controller集群元数据管理者从 Broker 中选出早期基于 ZooKeeper新版以内部 Raft 协议 KRaft 实现去掉了外部依赖执行故障转移第 4.2、4.3 两节合起来就是任何一台机器宕机只是丢了一个数据分片的某个副本其他机器上的副本可以立即顶上生产和消费完全不中断、数据不丢失的完整实现。其中 Unclean 选举是一个经典的可用性 vs 一致性开关金融级场景宁可短暂不可写也不允许落后副本当选丢数据。4.4 元数据管理集群还需要一个权威的元数据源哪些 Broker 存活、每个分区的 Leader 是谁、ISR 有哪些。Controller 扮演这个角色并广播变更配合 leader-epoch 机制保证副本恢复时日志截断的正确性避免脑裂场景下已提交数据被回滚。生产者/消费者客户端缓存元数据并监听刷新故障转移后自动把请求切到新 Leader。五、全链路数据不丢失ack 机制多副本解决的是存储冗余不丢失还需要写入确认协议贯穿生产端 → Broker → 消费端全链路。先看风险地图环节风险场景对策生产端 → Broker网络抖动消息没送达发送成功但 ACK 丢失ACK 确认 失败重试 幂等去重Broker 内部写了内存没刷盘就宕机多副本同步优先于单机刷盘宕的是一台不是整个 ISR副本间副本没同步完 Leader 就确认acksallISR 全体确认才返回 ACKBroker → 消费端消费者拿到消息处理一半宕机处理成功才提交 offset宕机后重投给其他实例5.1 生产端acks 确认 重试 幂等生产者写入的 ACK 级别参考 Kafka 的 acks 参数acks语义可靠性吞吐0发出去就算成功不等确认可能丢最高1Leader 写入本地即确认Leader 宕机且未同步时丢高all / -1ISR 全体写入才确认不丢配 min.insync.replicas中acksall时的完整确认协议时序两个必须配套的机制重试没收到 ACK 就重发设总超时上限 无限重试次数幂等重发必然带来ACK 丢了但消息已写入的重复。幂等生产者给每个 Producer 分配 PID每个 PID, 分区 维护单调序列号Broker 发现序列号重复直接丢弃——重试不再产生重复消息且不乱序。生产端推荐配置acksall # ISR 全体确认 retries2147483647 # 足够大的重试次数 delivery.timeout.ms120000 # 重试总超时上限 enable.idempotencetrue # 幂等生产者PID 序列号去重 max.in.flight.requests.per.connection5 # 幂等开启时仍保证分区有序5.2 Broker 端副本数与刷盘策略服务端的不丢失靠三条配置互相配合replication.factor3 # 每分区 3 副本分布在不同 Broker min.insync.replicas2 # ISR 至少 2 个才接受 acksall 写入 unclean.leader.election.enablefalse # 禁止落后副本当选 Leadermin.insync.replicas2与acksall组合的含义写入至少落在 2 个副本上才算成功此时同时挂掉 2 台机器数据依然不丢ISR 收缩到低于 2 时宁可报错拒绝写入也不静默降级——这是把不丢失置为最高优先级的显式表达。刷盘上选择依赖副本冗余而非单机 fsync异步刷盘 多副本确认。理由数据同时存在于多台机器的页缓存中单机掉电丢失的概率被副本数指数级稀释对金融级场景再叠加同步刷盘配置即可。5.3 消费端处理成功才提交宕机重投消费端参考 RabbitMQ 的全链路 ACK 思想消息从投递到确认之间存在生命周期未确认的消息在消费者宕机后必须重投。落到 Pull 模型上确认的载体就是 offset 提交关键纪律顺序不能反先处理后提交处理成功才提交 offset。反过来就是 At most once宕机丢消息重试上限 死信队列一条毒消息解析失败、依赖持续报错反复重试会卡住整个分区超过阈值进入死信队列DLQ人工介入不阻塞主流程消费幂等宕机重投必然带来重复消费业务侧用唯一键/去重表/状态机把至少一次收敛为效果一次。消费端推荐配置group.idorder-service enable.auto.commitfalse # 关闭自动提交处理成功后手动提交 auto.offset.resetearliest # 无位点时从头消费 max.poll.records200 # 单批拉取量保证处理时间不超过 poll 间隔 session.timeout.ms30000 # 心跳超时判定死亡5.4 消息生命周期全景把全链路串起来一条消息的状态机六、性能设计汇总前面各章的决策汇总成一张性能优化清单优化点手段原理收益顺序写Append-Only 追加日志消除磁盘寻道顺序写带宽百倍于随机写写入吞吐的数量级提升页缓存依赖 OS Page Cache 而非应用堆内存无 GC 压力、重启缓存仍在热数据读写近内存速度零拷贝sendfile 直接从页缓存到网卡消除用户态拷贝与上下文切换消费端 CPU 大幅下降稀疏索引 mmap每 4KB 一条索引项索引极小、整个驻留内存O(log n) 定位无随机 IO批量 压缩攒批linger.ms / batch.size lz4/zstd摊薄网络与 IO 固定开销吞吐成倍提升分区并行Topic 多 Partition 分布多 Broker读写并行度 分区数水平扩展无上限Reactor 网络模型少量线程 多路复用处理海量连接避免线程膨胀高并发连接低开销消费位点外部化Broker 不维护逐条消息状态服务端近乎无状态支撑百万级分区/消费者值得强调的取舍快大多来自减法——不维护复杂数据结构只追加、不做随机读写顺序扫描、不经过用户态零拷贝、不在服务端记状态位点外移。消息中间件把数据库随机读写 精确索引的能力换成了顺序读写 简单索引于是快了几个数量级。七、总结设计决策地图#问题我的设计业界参考1消息放哪内存页缓存缓冲 磁盘顺序追加 按大小/时间滚动的 SegmentKafka 文件存储模型2怎么定位消息分区内单调 offset 稀疏索引 mmapKafka offset / .index3怎么投递Pull 为主 长轮询兜底消费者组组内竞争、组间广播Kafka 消费模型 RabbitMQ 长轮询思想4怎么扩展Topic 分区打散到多 Broker并行度 分区数Kafka Partition5怎么高可用每分区多副本 ISR 动态收缩 Controller 故障转移选主Kafka 多副本 / ISR 机制6怎么不丢生产端acksall 无限重试 幂等去重Kafka 幂等生产者7怎么不丢Brokerreplication3 min.insync.replicas2 禁止 Unclean 选举Kafka 可靠性配置8怎么不丢消费端处理成功才提交 offset 宕机重投 死信队列 业务幂等RabbitMQ 全链路 ACK 思想9投递语义默认 At least once 幂等收敛为效果 Exactly once行业通用实践一句话总结这套设计用顺序 IO 与零拷贝换性能用分区与多副本换扩展性和可用性用贯穿生产端、Broker、消费端的全链路 ACK 换可靠性——三者共同支撑一个日均 TB 级、高性能、高可用、数据不丢失的消息中间件。延伸阅读建议研究源码的切入点Kafka《Kafka 权威指南》 源码LogAppend/ReplicaFetcherThread/GroupCoordinatorRabbitMQrabbit_queue/unacked消息状态机与镜像队列Quorum Queue实现对比 RocketMQ 的存储层CommitLog 单文件 ConsumeQueue 索引理解另一条设计路线的取舍。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询