手写Kafka级消息队列内核:日志存储、副本同步与性能哲学

📅 发布时间:2026/10/6 3:10:34
手写Kafka级消息队列内核:日志存储、副本同步与性能哲学
很多人觉得 Kafka 难懂是因为一上来就面对 Broker、Partition、Segment、ISR、控制器这一堆名词。我的学习路径可能比较笨既然看源码看不懂那我就自己动手写一个简化版的消息队列内核。这篇记录就是我当时从零实现Kafka 级消息队列内核的全过程——不是复刻完整的 Kafka那需要几十万行代码和大量的分布式协调工作而是把 Kafka 内核最核心的设计决策逐步落地日志的物理存储、生产者的批量缓冲、消费者的拉取与位移管理、高水位与副本同步以及零拷贝和页缓存背后的性能哲学。整个过程做下来最大的收获是明白了 Kafka 的很多反直觉设计其实都是同一套存储模型推导出来的。如果你也想深入理解消息队列原理或者是准备 Kafka 相关面试的后端开发者我建议你也走一遍这条路不用多一个能存、能取、能复制、能容错的最小内核就够用了。1. 先搞清楚Kafka 的内核到底指什么1.1 倒过来思考从提交日志出发我动手之前做了一件事把 Kafka 的所有功能列出来然后划掉外围的只保留内核。划完之后发现一个惊人的规律——Kafka 的本质其实是一套分布式的、分区的、带副本的提交日志Commit Log。换句话说所有高端功能都是从日志这个不变式上演生出来的。你可以把每个 Topic 的每个分区想象成一摞只能从尾部追加的作业本生产者的每条消息就是往当前这本作业本的末尾添一页写满一本就换一本新的这就是 Kafka 的日志分段Segment每个作业本都有一个编号编号从 0 开始单调递增这就是 Offset偏移量消费者做的事情更简单拿着自己上次读到的页码继续往后翻。我在设计内核时的第一行代码就是定义这个日志抽象。所有模块——网络层、复制机制、消费者位移——都必须围绕它构建。public interface Log { // 追加一条消息返回分配到的 offset long append(byte[] record) throws IOException; // 从指定 offset 开始读取最多读 maxBytes ByteBuffer read(long startOffset, int maxBytes); // 当前日志的末尾偏移量 long logEndOffset(); // 已经安全到达所有副本的偏移量 long highWatermark(); }这个接口看着简单但它把话说明白了内核的任务是存字节和取字节至于消息格式、压缩、权限、事务在外面包起来就好。1.2 最小内核的边界能砍的都砍掉我列了一张裁剪表把不需要的东西全砍了保留内核砍掉外围分区日志的追加写和分段存储流处理 APIKafka Streams生产者批量发送与重试Kafka Connect消费者拉取模式与位移提交Schema Registry主从副本与高水位可见性事务与 Exactly-Once 的高级封装基于 Page Cache 的读取权限控制 ACL做这个裁剪的意义不是偷懒而是明确边界我要验证的是 Kafka 的核心路径到底是怎么跑通的。如果连加了幂等、事务、多租户这些词都还没有理解日志模型那只会被复杂的功能淹没。1.3 立项前的技术选型为什么用 Java我用 Java 来实现理由有三一是 JVM 生态对缓冲区管理FileChannel、MappedByteBuffer和网络模型NIO支持好跟 Kafka 原生语言栈一致二是我后续可以直接对照 Kafka 源码逐行看它是怎么处理边界情况的三是做性能测试时Java 的 JMH 工具链成熟。如果你的目标是理解原理其实任何语言都能做——Go、Rust、C 都行。但别选那种写 IO 十分费劲的脚本语言因为消息队列内核一天到晚在跟字节打交道性能分析会很痛苦。2. 存储层从 Segment 到 Offset 的物理设计2.1 顺序追加与文件命名offset 即文件地址打开 Kafka 的数据目录你会看到每个分区一个文件夹里面是一堆.log文件和.index文件像这样topic-partition-0/ ├── 00000000000000000000.log ├── 00000000000000000000.index ├── 00000000000000036872.log ├── 00000000000000036872.index每个.log文件就是一个 Segment。文件名的数字是这个 Segment 里第一条消息的 offset。这是理解 offset 的关键offset 不是文件内部的行号而是全局单调递增的消息序号。消费者说我要从 offset 1024 开始读Broker 先根据文件名找到 00000000000000000000.log 不满足找到 00000000000000036872.log——好目标在新一段再从段内找。我的实现里段的滚动逻辑就一句话当前文件大小超过阈值我设置 128MBKafka 默认 1GB或者写满索引文件就创建新段。关闭当前段把文件名改成baseOffset.log新段 baseOffset 等于当前 LEO。public class LogSegment { private final FileChannel channel; private final long baseOffset; // 段内第一条消息的 offset private long nextOffset; // 下一个可用 offset private long sizeInBytes; public long append(byte[] record) throws IOException { ByteBuffer buffer encode(record, nextOffset); channel.write(buffer, sizeInBytes); sizeInBytes buffer.limit(); return nextOffset; } }为什么这么设计核心动机就一个顺序写。机械硬盘的顺序写和随机写差距可以到三个数量级SSD 虽然没那么夸张但顺序写依然比分页刷脏要稳得多。2.2 消息的物理格式长度前缀是万物通用语言消息落到磁盘上不能裸存一条字节数组需要携带控制信息。我设计的简化消息格式如下[4字节 总长度] [8字节 offset] [4字节 消息体长度] [N字节 消息体]总长度放在最外层的意义在于读日志时先读 4 字节拿到整条记录长度再按长度一次性读完整条消息。这种长度前缀Length-Prefix Framing在 TCP 分包、文件存储、网络协议里到处都是我之前总在代码里见到直到自己写日志文件才真正理解它存在的必要性——没有长度前缀你就永远不知道一截字节流从哪里断开是一条消息。真实 Kafka 的消息格式更复杂有批次Batch、CRC 校验、时间戳、压缩格式标记等。我的版本去掉了 CRC但在 offset 定位上保留了逻辑移位的思路。2.3 稀疏索引用可控的空间代价换 O(1) 定位光有.log文件还不够。如果消费者拿着 offset4096 来读而段文件是 128MB 的大文件你不能从头扫描吧所以每个段配一个稀疏索引文件.index。Kafka 的索引不是每条消息都建索引而是每写满一定字节默认log.index.interval.bytes4096才追加一条索引项。索引项记录的是[offset 的相对位移消息在文件中的物理位置]。class IndexEntry { long relativeOffset; // offset - baseOffset int position; // 物理位置 }查找流程是这样的根据目标 offset用二分法定位它所属的 Segment通过文件名 baseOffset加载该 Segment 的索引二分查找最后一个小于等于目标 offset的索引项从索引给出的物理位置开始顺序扫描直到读到目标消息。这里有个精妙的地方索引目录的二分查的是定位到大概位置最后用顺序扫描补齐剩余差距因为每 4KB 一条索引最坏情况也就是多扫描几 KB 数据。空间换时间换来的是接近 O(1) 的随机读。我当时测过一个数据在 128MB 的段上随机定位写在第 70MB 的消息加上二分和扫描总共耗时几百微秒级别。这个速度完全可以接受而且磁盘上索引文件占比只要很小一部分非常适合老数据长期保留的场景。3. 生产者链路批量、刷盘与持久化的取舍3.1 批量收集器吞吐的起点如果生产者一条一条往 Broker 发消息每条消息都要走一遍 TCP 握手和 syscall吞吐会非常难看。Kafka 用的是一个批量收集器Accumulator生产者先把消息攒在内存里攒够一批再一次性发出去。我在设计时很自然地想到了一个内存队列MapTopicPartition, DequeProducerBatch。每条消息进来时优先塞进当前分区未满的批次里批次满了就新建批次入队后台发送线程把批次序列化成网络请求。public class Accumulator { private final ConcurrentMapTopicPartition, DequeRecordBatch batches; private final long batchSize; // 批次大小阈值默认 16KB private final long lingerMs; // 等待时间阈值默认 0ms public void append(TopicPartition tp, byte[] record) { DequeRecordBatch deque batches.computeIfAbsent(tp, k - new LinkedList()); synchronized (deque) { RecordBatch tail deque.peekLast(); if (tail null || !tail.tryAppend(record)) { RecordBatch batch new RecordBatch(batchSize); batch.tryAppend(record); deque.addLast(batch); } } } }关键参数有两个batch.size控制攒多少才发linger.ms控制最多等多久。Kafka 默认linger.ms0意味着只要有发送线程空闲就立刻发不会故意等时间来凑批而batch.size到量后一定会触发发送。这就是为什么 Kafka 单分区吞吐可以非常高——每批次网络的往返成本被分摊到了很多条消息上。反直觉的点来了批量收集器并不是越攒越好它跟延迟存在天然矛盾。如果把linger.ms调到 50ms极端情况下每条消息都会多等 50ms消息延迟直接飙高。这正好回应了很多人问的Kafka 消息延迟高怎么排查——先去查是不是生产端把批量积累的时间调得太大。3.2 刷盘时机页缓存与 fsync 的暧昧关系写日志文件时有个永恒的纠结写入的数据什么时候真正落到磁盘Kafka 的处理跟传统数据库很不一样。传统数据库讲究提交后必须落盘所以事务提交要等 fsync保证掉电不丢数据。Kafka 默认却不这么做Broker 把消息写进操作系统的页缓存Page Cache就算写成功了fsync 的触发频率极低默认 Linux 上几乎是交给内核自己去脏页回写。我第一次看到这个设计很困惑那机器断电不是全丢了吗后来我意识到Kafka 的可靠性模型不是建立在单机落盘上而是建立在多副本上数据写入 Leader 的页缓存副本从 Leader 拉取数据也写入自己的页缓存只要acksall说明所有 ISR 副本都拿到了数据——这时候单机掉电其他副本还有完整数据fsync 只是为了应付整个集群全部掉电这种极端场景。我自己的内核实现里刷盘策略选了最简单的双阈值每写满 64MB 或者每 10 秒强制force()一次。这个策略在单机上可用但我很清楚它和 Kafka 的哲学不一样。Kafka 把持久性完全交给复制协议解决本地文件反而不着急落盘。如果你想在自己的项目里提高数据安全等级调 Kafka 的flush.messages和flush.ms能加速刷盘但副作用是吞吐下降这个权衡要自己把握。3.3 幂等生产者从源头减少重复既然消息可能因为网络重试导致 Broker 收到两遍Kafka 在服务端做了幂等每个生产者有个 producerId每条消息带一个自增的 sequenceBroker 端记录某个生产者最近写入的 sequence如果收到的 sequence 不是连续的就拒绝写入。我的简化实现也做了类似的事Broker 为每个 topic-partition 维护一个expectedSequence映射生产者批次里带baseSequenceBroker 校验后推进。这样虽然不能覆盖客户端崩溃重来的极端场景但足以解决同一条数据因网络瞬时抖动被发了两次的高频问题。4. 消费链路拉取模型、位移提交与重复消费的根4.1 为什么要学 Kafka 坚持拉而不是推很多消息中间件比如早期的 ActiveMQ、RabbitMQ 的某些模式是 Broker 主动推消息给消费者。Kafka 从一开始就选择了消费者主动拉取我当时也想过推不是更实时吗为什么逆潮流实现到一半我彻底理解了。推模型最大的痛点是背压问题消费者处理不过来怎么办Broker 要么给消费者加缓冲要么阻塞发送线程要么丢消息。消费者缓冲越大崩溃时重复消费的损失越大阻塞发送线程则直接影响其他分区。而拉取模型天然解决了背压消费者自己能吃多少就拉多少永远由消费者控制节奏。还有一个隐性好处拉取可以让消费者批量地拿数据。Broker 收到一个拉取请求可以把多条消息一次性返回网络 IO 效率高很多。推模型要做到这一点很别扭因为每个消费者处理速度不同Broker 顶多按固定速率推。// 消费者侧的核心请求 public FetchResponse fetch(FetchRequest request) { long offset request.getOffset(); LogSegment segment log.findSegment(offset); ByteBuffer data segment.read(offset, request.getMaxBytes()); return new FetchResponse(data, log.highWatermark()); }消费者拿到一批消息后只要维持住了拉了多少、处理到哪这两个状态Broker 完全可以无状态化处理消费请求。4.2 位移提交机制at least once 的根源聊到重复消费必须先建立一个关键认知Kafka 默认的消费语义是 At Least Once。为什么因为消费者先处理消息再提交位移offset。如果它在处理完消息后、提交位移前发生了崩溃下次重启时它会从旧位移继续读取那批消息就会被再读一遍。我实现的消费内核里位移管理长这样public class ConsumerCoordinator { // 记录每个分区“已经安全提交”的 offset private final MapTopicPartition, Long committedOffsets new ConcurrentHashMap(); // 记录本地已处理的 offset private final MapTopicPartition, Long processedOffsets new ConcurrentHashMap(); public void commitSync(TopicPartition tp) { committedOffsets.put(tp, processedOffsets.get(tp)); } public long nextFetchOffset(TopicPartition tp) { return committedOffsets.getOrDefault(tp, 0L); } }真实 Kafka 的位移不是存在消费者本地而是作为一条特殊消息写到内部 Topic__consumer_offsets由消费组协调器统一管理。但逻辑本质一样提交位移的时机 决定重复消费窗口大小的唯一变量。很多人面试被问怎么解决重复消费其实第一步要回答的是你在使用什么语义如果是 At Least Once那消息队列层面注定可能重复业务侧必须自己做幂等。4.3 重复消费的三大场景与应对手段我的内核实测中最容易出现重复消费的场景有三个场景一消费者崩溃后重启位移回退。处理了 100 条消息提交了 offset100然后还没来得及提交 offset200 就宕机了。重启后从 100 开始101~200 又消费一遍。这是最普遍的场景。场景二扩容/缩容导致的再平衡Rebalance。一个消费者组从 2 个消费者变 3 个某个分区要从消费者 A 转给消费者 B转让的时机是处理完当前批次但还没提交位移——于是 A 已经处理过的消息B 又消费一次。场景三网络闪断导致 Broker 端重试。消费者拉取数据后还没处理完TCP 连接断了Broker 认为消费者已死。重平衡后另一个消费者从旧位移继续拉——又重复了。应对手段其实没有银弹但我测试下来最有效的是三种组合手段原理适用场景业务侧幂等写入去重表/幂等键用唯一业务键去重重复消息被数据库拒绝大部分后台业务最推荐消费侧事务处理消息和提交位移放在同一事务里原子提交需要精确状态更新的场景成本高但语义最强至少一次 人工兜底接受重复靠监控和补偿任务清理脏数据日志分析等对账宽松的场景做这个内核让我彻底想明白一件事消息队列本身其实没办法做到完美的刚好一次最终还是要业务侧助一把力。Kafka 的 Exactly-Once 语义只在读 Kafka - 写 Kafka的流处理链路里有实现通过事务协调器换成取消息 - 写数据库就无能为力了。这个边界认知比背任何面试题都值钱。5. 副本与 ISR单机变集群的关键一步5.1 一个朴素问题机器挂了怎么办单机版的日志系统很简单但生产环境不可能让你揣着单点跑。我做副本功能时先思考了一个朴素问题Leader 收到消息写进了自己的日志Follower 怎么把数据追平Kafka 的方案是Follower 不主动接收数据而是向 Leader 发起拉取请求——Follower 把自己当前 LEO 传给 LeaderLeader 返回从 LEO 开始到最新的数据。这个设计和消费者拉取是一模一样的本质是存储节点的同步 一条消费日志的特殊消费者。// Follower 侧的复制线程 while (!shutdown) { FetchResponse response leaderClient.fetch(new FetchRequest(currentLeo)); for (Record record : response.records()) { localLog.append(record.getPayload()); } currentLeo localLog.logEndOffset(); saveReplicaOffset(currentLeo); Thread.sleep(replicaFetchMaxWaitMs); }这比我想象中简单太多一个 while 循环加追读就完成了主从同步。真正的复杂度不在于如何同步而在于同步到什么程度才算成功。5.2 高水位与可见性消费者看到什么才算安全引入副本后LEOLog End Offset和高水位HighWatermark就成了两个最关键的概念。LEO当前副本上日志的最后一条的偏移量1表示我本地写到哪了。HW高水位所有 ISR 副本都同步到的偏移量表示可以安全消费的位置。消费者能消费到的最大 offset 被限制在 HW 之内HW 以上的数据只有 Leader 自己能看到还没被确认复制到多数副本随时可能因 Leader 切换而丢失。我的实现中HW 的更新很简单Leader 定期比较所有跟随者的 LEO取最小的作为 HW。如果一个 Follower 落后太多被踢出 ISR它就不再参与 HW 计算。public long updateHighWatermark() { long minLeo leaderLeo; for (Replica r : isr) { minLeo Math.min(minLeo, r.logEndOffset()); } this.highWatermark Math.min(leaderLeo, minLeo); return highWatermark; }这个设计看着很简单但它回答了消费什么时候能看到数据的问题不是 Leader 一收到就可见而是等副本确认之后才可见。所以 HW 就是一条可视性红线它保证了消费者看到的数据在大多数副本上是存在的。5.3 副本追赶与 ISR 的伸缩规则ISRIn-Sync Replicas同步副本集合是 Kafka 巧妙地处理机器性能参差不齐的方案。不是所有副本都参与 HW 计算只有跟上节奏的副本才在 ISR 里。我的判出规则是Follower 超过 5 秒没向 Leader 发起拉取请求或者拉取位置上差距超过一个阈值Leader 就把它踢出 ISR。这样保证了一台慢机器不会拖垮整个集群的消费可见性代价是该副本暂时不参与 HW 的推进。这里有一个经典坑临时 GC 停顿导致的副本踢出。JVM 里一次 Full GC 可能超过 10 秒Follower 恰好在 GC 期间没法拉数据就被 Leader 踢出 ISR。踢出后Follower 追上进度会再次加入 ISR但期间集群的 HW 一度停止推进消息延迟肉眼可见。真实 Kafka 里这一步有追赶加回机制replica.log.end.offset追平后重新加入 ISR我的版本也实现了同样的逻辑并且我会在监控里盯UnderReplicatedPartitions指标任何超过几秒的不同步都是需要排查的信号。6. 性能底牌页缓存、零拷贝和批量为什么 Kafka 敢赌顺序 IO6.1 顺序 IO 的物理直觉与实测我先做了个简单的 IO 基准测试在普通 SSD 上往空文件里顺序写 512MB 数据再随机写同样大小的数据。结果顺序写在 1.5GB/s 量级随机写掉到几十 MB/s——差距接近两个数量级。这个差距的根源在寻址。磁盘的访问时间由寻道时间位置移动 旋转延迟 传输时间构成。顺序写让磁盘或 SSD 控制器可以大批量连续搬运随机写则每次都打断这个流水线。SSD 没有机械寻道但随机写也要付出擦写放大和映射更新的代价。Kafka 敢把性能赌在顺序 IO 上因为它知道日志文件的读写模式天然是顺序的——生产者追加写消费者也基本是按 offset 顺序读。这是存储模型与硬件特性的一次完美配合。6.2 页缓存不自己缓存把缓存交给内核Kafka 有个反直觉的设计它自己几乎不做堆内缓存。消息写入 Broker 是进 OS 页缓存读取时优先命中页缓存而不是搞一个 JVM 堆内的 redis。这样做的好处非常实际OS 的页缓存由内核管理淘汰策略成熟LRU 变体内存压力大时会回收冷页不会像自研缓存那样 OOM进程重启页缓存还在。Kafka 重启后消费旧数据不需要从磁盘重新读直接命中缓存冷热数据自动分层刚写入的热数据全在内存老数据被淘汰到磁盘。我在实现中一开始是自己维护一个 ByteBuffer 缓存后来遇到一个很明显的问题WindowsSubsystem 等等不其实是 JVM 堆被缓存挤爆我才发现自己又重复造了轮子。删掉自研缓存后内存使用率肉眼可见下降偶发 GC 也少了。让内核做它擅长的事应用层只做它必须做的事——这是 Kafka 给我上的一课。6.3 sendfile 零拷贝省掉拷贝的本质传统读文件发送到 socket 的路径是这样的磁盘 - 页缓存DMA 拷贝页缓存 - 用户缓冲CPU 拷贝用户缓冲 - socket 缓冲CPU 拷贝socket 缓冲 - 网卡DMA 拷贝。整个过程四次数据拷贝其中两次 CPU 拷贝。如果数据量是 100GB 的日志扫描这两次 CPU 拷贝的耗时会非常可观。Kafka 使用sendfileJava 里对应FileChannel.transferTo把第 2、3 步直接交给内核完成数据从页缓存直接被 DMA 引擎送到网卡省掉了用户态的参与CPU 拷贝减到零或极低。这就是零拷贝。// 消费读取路径直接把 segment 文件内容送到 socket long transferred fileChannel.transferTo(position, count, socketChannel);没有什么花哨的就是一行调用。但你必须意识到它省掉的是什么当并发消费者多、读取数据量大时省掉的 CPU 拷贝时间能用多核每秒几万次系统调用的代价来衡量。我对比过普通 read write和transferTo两种路径处理 1GB 消息前者耗时高了不少后者几乎打满网络带宽。6.4 端到端压缩CPU 换 IO划算最后一个性能大招是压缩。Kafka 把压缩放在生产者客户端做Broker 和消费者在有需要时才解压。这有点反直觉但逻辑很通压缩需要 CPU但省下的是网络带宽和磁盘空间。我实测过用 gzip 压缩 JSON 数据压缩率能到 4 倍以上。代价是消费者要解压CPU 占用多了几个百分点。这就是典型的用 CPU 换 IO在网络和磁盘是瓶颈的场景里非常划算。如果你的消息本身就是压缩过的视频、图片字节流那再压就没意义可以跳过但常规文本日志、订单数据这类我推荐保持默认压缩配置。7. 从能跑到像 Kafka差距、测试与踩坑记录7.1 我踩过的三个坑坑一offset 错位导致读了半天全是脏数据。现象是连续写 10 万条消息重启后从 offset 50000 读读到的消息体错乱。排查到根因是我的消息格式里总长度字段和offset字段顺序写反了导致物理位置和逻辑偏移量对不上。这让我明白日志文件的格式设计是不能随意变更的协议一旦写入文件必须考虑向前兼容。后来我在格式头加了版本号字段。坑二同一分区的并发写导致日志顺序错乱。我的生产者客户端是多线程模型早期版本直接让多个线程同时调Log.append没有锁。结果落盘的消息前后顺序和 offset 分配顺序对不上。Kafka 的解法是同一个分区同一时刻只有一个生产线程能写入单写者模型靠锁保证。我把生产者端的批次分桶 - 单队列单线程发送做了之后问题消除。这也是 Kafka 单分区吞吐高而不是靠多线程写出来的原因。坑三消费者提交位移的竞态。我最早版本的位移提交是每消费一条就提交一次结果消费者线程和提交线程之间出现竞态两条线程同时提交后提交的旧位移覆盖了新位移。后来改成定时批量提交 提交前校对 processedOffset committedOffset才稳住了。这个改动直接让我的消费端从偶尔重复大量数据变成重复窗口可控。7.2 验证方法论光能跑远远不够做消息队列你怕的不是程序崩溃而是数据静默丢失。所以我给自己定了一个验收清单顺序校验写入端为每条消息附上序号消费端校验序号是否严格递增一旦发现跳号或乱序立刻报错重复率观测消费者记录每个 offset 被消费的次数统计重复率。起步阶段我的重复率高达 5%把位移提交改成批量提交后降到 0.2%副本追赶时间挂掉一个 Follower恢复后记录它追上 Leader 的耗时验证 ISR 重加入逻辑Kill 测试随机 kill 消费者、kill Broker、kill 网络看数据是否丢失、重复是否可控。这些测试不需要复杂框架写几个脚本就行。但价值很大——它们逼着我把看起来能跑变成在各种失败下都能跑。7.3 和真 Kafka 的差距对照与后续学习路线做完这个 mini 内核我给自己列了一张差距表也算是继续深入的方向维度我的实现真实 Kafka网络协议自定义简化协议自研二进制协议支持请求多路复用位移存储本地文件内部 Topic __consumer_offsets跨节点协调分区分配静态指定由 Controller 选举动态分配并管理副本清单事务与幂等仅幂等完整事务协调跨分区原子提交控制器无Controller 负责分区领导权分配与故障转移如果看完这篇你也想动手我建议按这样的顺序啃源码先读core/src/main/scala/kafka/log/Log.scala日志段的滚切与索引再读kafka/server/KafkaApis.scala请求处理最后读kafka/coordinator/group消费组协调。代码虽然 Scala 写起来有点门槛但核心逻辑跟我这个简化版几乎一一对应。最后说一个做完整套实现的体会消息队列这个领域面试题里的为什么基本都能在存储模型里找到答案。Kafka 的高吞吐源于顺序 IO 和批量高可用源于副本和 HW重复消费源于位移提交的时机数据不丢源于 acks 与 ISR 的配合——这些东西不是独立的技巧而是一棵日志树上长出来的不同分支。写完这个内核之后我再去看任何消息队列的架构文档感觉完全不一样了不是在看一个新系统而是在看同一个模型在不同约束下的变体。如果你正在被 Kafka 的原理绕晕不妨也试着从日志这个原点出发按自己理解写一遍收获会远超预期。