
聊到消息队列大家第一反应基本都是 Kafka、RocketMQ、RabbitMQ 这些成熟产品。但我今天想分享的是一个有点“重造轮子”味道的经历从零实现一个 Kafka 级消息队列内核。先解释一下这里的“内核”不是 Linux 内核而是消息队列 Broker 的核心程序相当于 Kafka 的 Server 端最底层那套存储、网络、复制与消费协调逻辑。完整复刻 Kafka 全部特性不现实也不必要但把核心子集吃透你能回答“Kafka 为什么这样设计”“offset 为什么是那个语义”“重复消费到底卡在哪”这些面试和实战里都要命的问题。这篇博文适合三类人想深入理解消息中间件原理的后端工程师、需要在低资源或定制场景里做轻量消息系统的架构师以及准备消息队列相关面试的人。我会把整条实现路径拆成六个部分整体设计、存储引擎、网络协议与线程模型、副本与消费组、选型对比与实战建议、常见问题排查实录。每一步都会用实际可抄作业的伪代码、参数表和排查清单来说明而不是停留在概念层面。1. 整体设计先想清楚一个消息队列内核到底要解决什么1.1 这不是重新发明轮子而是为了回答“Kafka 为什么这么设计”很多人觉得自研消息队列不划算我同意大多数生产场景直接用 Kafka 就好。但如果你遇到下面这几种情况“从零实现”就有它不可替代的价值你需要把核心原理彻底搞懂而不只是会调用 API。你的业务对消息中间件有极强的定制需求比如要内嵌到嵌入式环境、要精简到只有几 MB 内存、要控制每个文件格式。面试中经常被问到 Kafka 源码级问题但只看源码容易迷失动手做一遍能把“线性读”“页缓存”“零拷贝”“ISR”“高水位”这些概念全部串起来。这个项目定位是“Kafka 级”但不等同于“Kafka 完整版”。我会实现它的核心子集高吞吐持久化、分区有序、多消费者组、副本同步、消费再平衡。像权限认证、多租户、配额管理、Exactly-Once 这种外围能力可以在核心跑通之后再慢慢加。1.2 核心抽象从 Broker、Topic、Partition 到 Offset动手之前先要把消息队列的抽象模型定下来。Kafka 的模型是这样Broker一个独立进程负责存储和转发消息。Topic业务上的逻辑分类比如“订单事件”“用户行为日志”。PartitionTopic 下的物理分片是并发和扩展的最小单位也是消息顺序的边界。OffsetPartition 内单调递增的消息序号消费进度就是通过这个整数记录的。Consumer Group组内消费者共同消费一个 Topic每条消息只会被组内一个消费者处理不同组之间则各自独立收到所有消息。这套模型的精髓在于“顺序降维”。如果要求整个 Topic 全局有序吞吐量会非常难看因为全局锁会成为最大瓶颈。Kafka 把“顺序”限制在单个 Partition 内Broker 层面的并行度就等于 Partition 数量乘以副本数。自研内核时我也直接沿用这个设计避免一开始就掉进“既要全局有序、又要高吞吐”的坑。1.3 整体模块划分与关键技术选型我的自研内核拆成了五个模块模块职责自研方案存储模块消息落盘、索引、日志清理append-only 文件 稀疏索引不依赖数据库网络模块客户端接入、请求分发、长轮询Netty 实现主从 Reactor 模型复制模块Leader 与 Follower 同步、故障转移ISR 机制 高水位提交语义消费协调模块消费组管理、分区分配、再平衡组协调器模型管理模块Topic 管理、配置、监控RESTful API JMX 指标为什么不直接拿 MySQL 存消息最主要原因是写入模型差太远。数据库写入要维护 B 树索引、undo log、redo log还要处理事务锁和并发控制随机写和写放大非常严重。而 Kafka 能扛住百万级 TPS核心是把消息当成“日志文件”直接 append 到文件尾部不修改已有数据配合操作系统的页缓存完成写入加速。自研内核时存储这条路必须走 append-only。2. 存储引擎消息队列吞吐的灵魂2.1 append-only 日志写入为什么这么快现代 SSD 顺序写可以轻松达到几百 MB/s随机写却经常掉一个数量级。append-only 的核心就是只往后追加不更新旧数据所以磁盘写入非常友好。实现上一个 Partition 会对应一个目录目录下有一组 segment 文件。每个 segment 是固定大小上限的日志文件比如 1GB。写入时先定位到当前活跃 segment 的尾部把消息按约定格式追加进去。这里有一个关键格式RecordBatch。Kafka 不是一条一条写消息的而是攒成一个批量一起写。在自研内核里我定义了这样的字节布局字段长度说明baseOffset8 字节批次内第一条消息的 offsetbatchLength4 字节整个批次字节数partitionLeaderEpoch4 字节分区 leader 的 epoch防止脑裂写入magic1 字节消息格式版本用于兼容升级crc4 字节批次校验值attributes2 字节压缩算法、时间戳类型等标记lastOffsetDelta4 字节批次内最后一条消息与 baseOffset 的差值records变长具体消息列表写入伪代码看起来就是这样public AppendResult append(RecordBatch batch) { FileChannel channel currentSegment.channel(); long baseOffset currentSegment.nextOffset(); ByteBuffer buf encodeBatch(batch, baseOffset); // 重点定位到文件尾再写避免覆盖已有数据 channel.position(channel.size()); while (buf.hasRemaining()) { channel.write(buf); } // 可选刷盘见 2.3 节 maybeFlush(); currentSegment.incrementNextOffset(batch.count()); }我在实现中踩过的一个坑是字节缓冲区的复用。刚开始我给每个写入请求新建 ByteBuffer结果大量消息时 GC 压力很大。后来改成线程私有缓冲区写入高峰期表现稳定很多。另外文件预分配也很有用提前把文件扩展到预定大小避免写入时频繁触发文件系统元数据更新。2.2 索引文件offset 与时间戳的快速定位光有日志文件还不够消费者要按 offset 拉消息不能每次从头扫描整个文件。所以每个 segment 除了日志文件还要配套两个索引文件offset 索引和时间索引。offset 索引的结构是一个稀疏索引并不是每条消息都记录。我设置的策略是每隔 4KB 日志量写入一个索引项记录“相对 offset”和“物理位置”。相对 offset 是相对于 segment 起始 offset 的差值用 4 字节存储物理位置也用 4 字节这样一个索引项只有 8 字节1GB 的 segment 索引内存占用还在可控范围内。为什么用稀疏索引而不用全量索引因为顺序日志文件里相邻消息的物理位置也是相邻的。消费者要拉 offset 为 12345 的消息先在索引文件里二分查找小于等于 12345 的最大索引项拿到一个物理位置然后从那个位置开始顺序扫描日志通常读几十字节就能定位到目标消息。这种“索引定位 局部扫描”的策略用极小的内存换来了可接受的定位延迟。时间索引的原理也类似只不过 key 是消息时间戳value 是相对 offset。这个索引主要用于“从某个时间点开始消费”的场景。Broker 重启恢复时不需要扫描全量日志只需要从最近一次 checkpoint 开始对活跃 segment 的末尾做一次扫描重建缺失的索引项就行。2.3 页缓存、零拷贝与操作系统合作一个容易忽略的事实是Kafka 的写入并不强制每次都 fsync 到磁盘。Producer 发过来的消息先进入内核页缓存page cache由操作系统决定什么时候真正写回磁盘。只要消费者读的速度跟得上消息直接命中页缓存根本不会产生磁盘 IO。自研内核也沿用了这个设计写入路径不调用 force()只依赖 OS 的 pdflush 机制。读路径上我做了两件事第一消费者拉取日志时使用FileChannel.transferTo()方法来实现零拷贝。它相当于系统调用 sendfile内核直接把页缓存中的文件数据发送到 socket不需要经过用户态拷贝和应用层再封装。我在测试环境对比过普通 read 加 write 的方式大约只能跑零点几的带宽而 transferTo 可以直接打满机器网络带宽。原因就在于省掉了一次用户态到内核态的来回拷贝还有一次应用层缓存复制。第二页缓存命中率决定了读性能的上限。所以我会严格控制 segment 文件数量防止老 segment 频繁被换出缓存也会监控page cache hit ratio这个指标。如果命中率持续低于 90%意味着消费者的拉取速度跟不上生产速度需要扩容消费者或减少保留时长。刷盘策略必须做成可配置的。我提供了两个关键参数参数推荐值说明flushIntervalMs100每 100ms 强制刷盘一次flushMessages10000积压消息达到 1 万条时刷盘这样即使进程崩溃最多丢失几百毫秒的数据如果业务要求绝不能丢可以把刷盘间隔设为 0代价是吞吐明显下降。真实场景里“完全可靠”和“高吞吐”是鱼与熊掌必须让调用方有选择空间。2.4 日志清理不删除文件和删除无效数据是两回事消息队列不能无限存数据。日志清理有两种策略基于时间/大小的保留retention和基于 key 的压缩compaction。保留策略很简单后台线程定期扫描 segment 的修改时间或文件大小删除过期 segment。但要注意删除操作不能影响正在读取该 segment 的消费者。我采用 F1 策略只删除文件句柄让已经打开这个文件的消费者继续读完等文件引用归零后再真正删除。这个细节在文件系统层面很关键。压缩策略就复杂一些需要保留每个 key 的最新一条消息因此 Kafka 也提供了后台 Cleaner 线程。自研内核可以先不做完整 compaction但至少要知道它的存在在一个消息携带 key 的业务场景里只保留最新值能大幅度降低磁盘占用。如果项目精力有限我建议把 compaction 放到第二阶段第一阶段的重点永远是 retention 删除否则磁盘会先爆。3. 网络协议与线程模型把吞吐量从磁盘搬到客户端3.1 自定义协议长度前缀、魔数与版本兼容存储再快网络层跟不上也是白搭。自研消息队列可以不用 Kafka 的二进制协议但协议设计必须考虑到三个问题粘包半包、版本兼容和请求响应匹配。我采用的协议帧格式是字段长度说明magic2 字节魔数标识协议类型version1 字节协议版本号type1 字节请求类型如 PRODUCE、FETCH、COMMIT_OFFSETcorrelationId4 字节请求与响应对应关系payloadLength4 字节后面 payload 的字节数payload变长具体业务参数或数据在 Netty 的 pipeline 里我用LengthFieldBasedFrameDecoder先解决粘包半包然后自定义的 Handler 负责解析 payload。每个响应必须带相同的 correlationId。这个字段看起来不起眼但排查线上问题的时候有它才能把客户端堆栈、Broker 日志和网络抓包对齐到同一条请求。早期版本我为了省 4 个字节没加后来排查超时问题差点被逼疯最后老老实实加回去。为什么不用 HTTPHTTP 本身可以完成请求响应但头部解析和序列化开销大而且长连接管理、流式响应这些能力对消息队列来说不够轻量。自定义二进制协议虽然在开发上麻烦一点但性能和可控性好得多。3.2 Reactor 线程模型与背压处理网络线程模型我用的是主从 Reactor这是 Netty 的标准模型主 Reactor 只负责 accept 新连接分配一个 Channel 给从 Reactor。从 Reactor 负责 socket 的读写解析出完整的请求帧后投递到业务线程池。业务线程池里的线程负责执行 append 和 fetch 等实际逻辑。这里最容易犯的错是让 IO 线程直接执行业务逻辑比如把文件写入放在 Netty 的 EventLoop 里。如果写盘变慢所有连接都会被拖垮。正确的做法是 IO 线程和业务线程彻底分离。背压同样不能忽略。消费者处理慢时Broker 不能无限往这个连接的内存队列里塞数据。我给每个连接维护了一个发送队列并设置一个高水位阈值。当队列积压超过阈值时暂停从该连接读取新请求让 TCP 缓冲区自然填满从而让消费者客户端感知到背压。这么做比在大内存里硬扛要安全得多因为堆内存被打满会导致 Full GC而这个连锁反应会把整个 Broker 拖停。线程数也不能盲目调大。我测试过业务线程池在 4 到 8 个线程时收益最明显超过 8 个以后上下文切换和锁竞争带来的开销会抵消并发收益。具体数字取决于磁盘类型和 IO 调度但核心思想是不要用大量线程去掩盖单线程 IO 上的瓶颈先用工具定位瓶颈再调参。3.3 消费者拉取模型Fetch 请求、长轮询和速率控制Kafka 选择消费者主动拉取模型而不是 Broker 推模型是有原因的。拉取模型里消费者根据自己的处理能力决定拉多少条天然具备流控能力Broker 也不需要在内存里维护大量推送连接。自研内核继续走拉取模型。当分区里暂时没有新消息时Fetch 请求会怎样如果立刻返回空结果消费者只能疯狂轮询浪费 CPU 和网络。我的做法是让 Fetch 请求在 Broker 侧挂起一段时间默认 500ms这期间有数据到达就立即唤醒返回如果超时还没有新消息就返回一个空批次。这种机制叫长轮询。通过这个设计消费端的空轮询频率可以降到每 500ms 一次同时消息延迟几乎不受影响。消费者拉取时要考虑三个参数minBytes至少要攒够多少字节才返回避免小批次过多。maxBytes单次 Fetch 最多返回多少字节防止一次响应过大。maxWaitMs最多等待多长时间控制延迟上限。默认配置下minBytes是 1maxBytes是 50MBmaxWaitMs是 500。如果业务需要低延迟就把maxWaitMs调小如果追求吞吐就调大minBytes让 Broker 多攒一会儿批量返回。4. 副本、ISR 与消费组从单机走向集群4.1 副本与 Leader 选举为什么不能只在单机单机消息队列只要进程一挂消息就全没了。要真正安全必须引入副本机制。我的设计里每个 Partition 的副本数可以配置比如 3。副本从逻辑上分为 Leader 和 FollowerLeader 负责处理生产者和消费者的请求Follower 只负责从 Leader 拉取日志并追赶进度。这里有两个概念绕不开LEO 和 HW。LEO 是 Log End Offset指的是日志最后一条消息的下一个位置。HW 是 High Watermark表示“已提交”位置它是所有 ISR 副本都复制到的位置。消费者只能读到 HW 之前的消息HW 之后的消息虽然已经在 Leader 上但还不算“已提交”。为什么要有 HW考虑这个场景Leader 收到消息后还没来得及给 Follower 复制机器就宕机了。如果新选出的 Leader 日志比旧 Leader 短消费者可能读到旧 Leader 又突然回滚的数据这在消息语义上是不允许的。HW 的作用就是定义“提交点”只有 ISR 全部确认过的位置才允许被消费。生产者端用acks参数控制数据的可靠性acks 值行为可靠性0发出去就完事不等确认最高吞吐可能丢消息1Leader 写入本地日志后确认最常见配置Leader 挂掉可能丢-1all等待所有 ISR 副本都确认可靠性最高延迟最高如果设置了acksall但 ISR 里只有 Leader 一个那其实也退化成 acks1。所以要同时配置min.insync.replicas2意思是至少两个副本确认才算成功否则返回异常给生产者。这套机制保证了单台 Broker 宕机时消息不丢。4.2 消费组再平衡重复消费问题的主要源头消费者组让一组消费者共同消费一个 Topic核心是协调器Coordinator分配分区的归属。消费者启动后会向组协调器上报自己加入组协调器根据当前成员把分区重新分配。一切看起来挺顺但正是“再平衡Rebalance”成了重复消费问题的主要源头。再平衡的触发条件有四个消费组成员发生变化比如新增、退出、宕机。订阅的 Topic 发生变化。分区数发生变化。消费者心跳超时被协调器判定为死亡。再平衡为什么会导致重复消费原因在于 offset 提交和业务处理的先后顺序。如果消费者先处理完消息还没来得及提交 offset这时候就发生了再平衡协调器会把分区分配给另一个消费者。新消费者从旧的已提交 offset 位置开始拉取于是之前已经处理完的消息会被再次消费。排查重复消费问题我有一套固定流程先看代码里 offset 提交的时机。如果是在finally块里提交并且是先处理业务再提交重复消费的概率高。检查消费者参数session.timeout.ms和max.poll.interval.ms。如果处理一条消息耗时太长超过了max.poll.interval.ms消费者会被判定死亡强制触发再平衡。在消费端做幂等处理。不管底层怎样重复业务层要保证同一业务标识重复执行不会产生副作用。幂等的实现方式包括数据库唯一索引、Redis 的 SETNX、业务自带的去重表、或者利用消息里的业务主键做过滤。这是最后一道兜底防线也是生产环境里的硬性要求。4.3 消息顺序性单分区有序的局限与多线程保证Kafka 的顺序性只存在于单分区内部跨分区之间不保证任何顺序。自研内核同样如此。如果业务要求严格的全局顺序最直接的办法是让这个 Topic 只有一个分区但这样并行度就没了。如果业务允许按 key 有序更合理的方案是按 key 做哈希将同一 key 的所有消息路由到同一个分区。消费端单分区内保持顺序处理或者按 key 分发到多个内部队列每个 key 固定由一个线程处理。我常用一个叫“分桶 单线程 worker”的模式public class OrderedConsumer { private final MapString, BlockingQueueMessage buckets new ConcurrentHashMap(); private final MapString, Thread workers new ConcurrentHashMap(); public void onMessage(Message msg) { // 相同 key 落同一个桶 String key msg.getKey(); BlockingQueueMessage queue buckets.computeIfAbsent(key, k - new LinkedBlockingQueue()); queue.put(msg); // 每个桶绑定一个专用线程串行处理桶内消息 workers.computeIfAbsent(key, k - { Thread t new Thread(() - { while (true) { Message m queue.take(); process(m); } }); t.start(); return t; }); } }这个方案的精髓是并发处理不同 key 的消息同时保证相同 key 的消息串行执行。它比单分区全局串行好因为不同 key 之间互不阻塞。代价是线程数可能很多所以桶数一般要做上限超过上限的 key 落到同一个兜底桶里。5. 选型对比与实战建议什么时候应该用成熟 MQ什么时候适合自研5.1 Kafka、RocketMQ、RabbitMQ 的取舍速查很多人纠结消息队列选型我也被问过无数次。这里给出一个实战向的对比表维度KafkaRocketMQRabbitMQ自研消息队列消息模型Topic 分区模型Topic 队列模型Exchange 路由模型自定模型吞吐量极高高中取决于实现延迟低极低低可控顺序消息分区内有序支持局部顺序需要单队列分区内有序事务消息支持原生支持不完整支持基本没有路由能力弱中强自定义运维复杂度高高低极高自研成本我的建议是日志采集、大数据管道、流处理场景首选 Kafka。金融交易、对事务消息和低延迟有强需求首选 RocketMQ。业务逻辑复杂、路由规则多、小规模场景RabbitMQ 更灵活。只有那些需要极强定制、低资源内嵌、或者纯粹为了学习内核原理的场景才适合走自研这条路。生产环境没必要硬刚自研但理解内核能帮你做更准确的选型。比如你看到某业务经常“消息延迟高”如果懂 Kafka 的 Fetch 长轮询机制和 Page Cache 原理你就知道先去看消费者处理耗时和 Broker 页缓存命中率而不是盲目调大线程数。5.2 从内核实现到生产可用的距离如果自研消息队列想真正上线还需要补齐很多周边能力元数据管理Topic、Partition、配置的持久化和一致性。权限认证至少要做 SASL 或 TLS 认证。配额控制限制单个客户端的生产消费速率防止恶意打爆 Broker。监控告警Broker 指标、消费延迟、网络吞吐量。协议兼容如果能让客户端直接使用 Kafka 原生协议就能直接用现成的 Kafka 可视化工具来观察 topic 和 offset省掉做管理界面的成本。我的建议是从小步快跑路线开始先跑通单节点再实现主从复制最后加上消费组协调。不要一开始就追求协议兼容 Kafka会陷入和官方协议定义的死磕而核心存储和网络还没验证。如果你的接口是 RESTful 风格的可以用最简单的 HTTP 客户端做测试等核心功能稳定后再考虑兼容原生协议。6. 常见问题与排查技巧实录6.1 消费者重复消费排查实录真实场景某团队使用自研消息队列处理交易事件偶发出现同一笔订单被处理两次。排查过程我记录如下第一步查消费端日志发现两次处理中间隔着一次 Rebalance日志里能看到“rebalance started”和“partition revoked”。这基本就锁定是再平衡导致的重复。第二步看 offset 提交代码发现处理器在消息处理之前就把 offset 提交了。这样就等于告诉 Broker “这条消息已经处理完了”但实际业务还没执行完毕。如果真的就那一刻挂了消息会直接丢失而如果赶上再平衡消息会被重新拉取。第三步修复策略改成两段式先处理业务再提交 offset同时把max.poll.interval.ms调大到业务处理耗时的 3 倍以上。第四步加幂等兜底。即使前面还有漏洞订单处理接口以“订单号”做唯一键重复请求直接返回之前的处理结果。最终结果重复率从千分之三降到了零。下面是一个检查点清单检查点现象对策offset 提交时机提交先于业务执行改为业务完成后提交处理超时消费线程处理超过 max.poll.interval调大超时或拆分逻辑Rebalance 日志日志中出现多次 rebalance检查心跳、网络抖动业务幂等重复数据影响最终结果数据库唯一键或去重表6.2 消息延迟高从不正常到定位“消息延迟高”是消息队列最常见的问题但原因五花八门。我一般从四个方向排查第一是 Broker 负载。用top和iostat看 CPU 和磁盘 IO 是否打满。如果磁盘 IO 利用率长期超过 90%说明日志刷盘或者消费者读日志的时候页缓存没命中读到了真实磁盘。这时候要看是不是段文件太多、页缓存命中率低。第二是网络带宽。sar -n DEV能看出网卡吞吐是否接近上限。消息队列的高吞吐场景经常先卡在网络而不是磁盘特别是压缩后带宽没有优化的时候。第三是消费者处理速度。如果消费者线程里有个慢 SQL消费速度低于生产速度Lag 会持续上涨。这时候延迟高其实不是消息队列的问题而是消费者自身的问题。第四是 Fetch 参数配置。如果maxWaitMs设置成了 1000ms而业务期望 100ms 内收到消息那延迟就锁死在 1 秒了。低延迟场景要把maxWaitMs调小同时配合minBytes让 Broker 尽快返回。6.3 关键自检清单与参数速查下面是我在自研内核里比较关键的参数也贴近 Kafka 的语义供大家对照参考参数含义影响log.segment.bytes单个 segment 最大字节数影响索引文件大小和日志清理粒度log.retention.hours消息保留时长影响磁盘占用和可回溯范围flush.interval.ms刷盘间隔影响吞吐量和故障丢失窗口num.io.threads业务处理线程数影响并发写入和读取能力fetch.wait.max.msFetch 请求最大等待时间控制消费延迟min.insync.replicas最小同步副本数控制数据可靠性下限如果你要复现整套实现建议用一个 3 节点集群每个节点上跑一个 Broker 进程磁盘用 SSD网络至少千兆。测试时先压生产者吞吐确认写入不丢消息再压消费者吞吐确认消费延迟最后做一次 Kill -9 测试观察数据恢复和重复消费是否在预期范围内。我自己完整走完一遍之后最大的感受是Kafka 的很多设计看似简单但每一个都是在极端场景下被逼出来的选择。比如 append-only 日志、页缓存、零拷贝这个组合缺一环都到不了所谓“Kafka 级”的吞吐ISR 和 HW 这套机制少一个细节都可能出现消息回滚或丢失。这也是我强烈建议所有后端工程师哪怕日常只用现成消息队列也值得动手写一个最小实现的原因。源码看得再多都不如自己踩一遍坑记性深刻。