从零构建 Kafka 后发现:消息生产只是顺序追加日志

从零构建 Kafka 副本后,作者确认消息生产时磁盘操作只是顺序追加到分区日志文件,而非随机写入。这直接解释了为什么分区必须有单一 leader——所有副本同时接受写入会破坏一致性。

作者用 Kafka 多年,熟悉 topics、partitions、producers、consumers 和 replication 这些术语,却始终觉得只知道名字,没摸到它们为什么这样组合。于是他决定动手写一个极简的消息代理,不是为了替代 Kafka,而是为了看清抽象下面的真实行为。

这个项目直接回答了他心中的几个核心疑问:消息落到磁盘时到底发生了什么?分区为什么非要选一个 leader?所有副本都能接收写入会怎样?消费者组加入新成员时如何分配分区而不重复消费?通过实际编码,这些问题不再停留在文档层面,而是变成了代码里能触碰的逻辑。

消息生产时磁盘上实际只做顺序日志追加

在自建版本里,生产一条消息的磁盘操作极其简单:把消息按顺序追加到对应分区的日志文件末尾。作者实现的日志模块就是一个不断 append-only 的文件,没有任何随机写,也没有复杂的索引更新操作。

具体过程是,producer 把消息发给 leader,leader 先把消息写入自己的日志文件,同时记录 offset。整个写入采用顺序追加方式,利用操作系统页缓存,速度远高于随机 IO。这解释了 Kafka 为什么能支撑高吞吐:磁盘顺序写几乎和内存写一样快。

自建代码里,日志文件以二进制格式存储,每条消息前面带长度字段,便于后续读取。作者刻意没有实现官方 Kafka 的零拷贝和复杂压缩,只是保留了最核心的 append 行为。测试显示,即使在普通机械硬盘上,顺序追加也能轻松达到每秒数万条。

这个发现让作者明白,Kafka 的高性能不是靠什么魔法算法,而是把磁盘最擅长的顺序写用到了极致。所有上层设计——分区、复制、索引——都是围绕这个“只追加”前提展开的。

分区 leader 存在是为了避免多副本同时写入冲突

为什么分区必须有一个 leader,而不能让所有副本都接受生产请求?自建过程中作者很快遇到了这个问题。如果允许任意副本接收写入,那么不同副本收到的消息顺序可能不一致,后续复制时就无法对齐。

在简化实现中,作者只让 leader 接受生产请求,其他 follower 只负责从 leader 拉取日志并追加到自己文件。leader 负责分配 offset,保证同一分区内消息全局有序。follower 复制时采用 pull 模式,定期请求 leader 上最新的 offset 范围。

如果去掉 leader,让所有副本同时写,会立刻出现脑裂:两个副本各自收到不同消息,offset 冲突,后续消费者看到的数据不一致。leader 机制本质上是把写入串行化,确保整个分区只有一条日志链。

这个设计代价是 leader 所在 broker 压力更大,但换来了强一致性。作者在代码里用简单的主从选举模拟了这个过程,虽然没有 ZooKeeper 或 KRaft 那么复杂,却足以让人看清 leader 的必要性。

副本复制在简化实现里如何取舍一致性与可用性

官方 Kafka 支持可配置的 replica 数量和 min.insync.replicas 参数,平衡一致性和可用性。自建版本中作者做了更极端的简化:leader 写完后立即通知 follower 拉取,follower 追加完成后才返回 ack。

这种同步复制保证了只要 leader 确认返回,数据就在多数副本上落盘,但也意味着生产延迟会随副本数线性增加。作者测试了 3 个副本的场景,发现延迟比单副本高出约 2 倍。

与官方实现相比,自建版缺少 ISR(In-Sync Replicas)动态维护机制,也没有 leader 选举时的 epoch 防脑裂。官方 Kafka 用 controller 和 ZooKeeper 精确管理哪些副本还在同步,自建版则简单假设所有 follower 都在线。

这些取舍让作者明白,生产级 Kafka 的复杂性很大一部分来自对故障场景的严谨处理。简化版虽然功能残缺,却把“复制就是把日志从 leader 拷贝到 follower”这个本质暴露得清清楚楚。

消费者组协调通过什么方式分配分区

消费者组是 Kafka 实现负载均衡和 exactly-once 语义的关键。自建版本中,作者实现了一个极简的组协调器。新消费者加入时,协调器会把当前所有分区重新分配给组内所有消费者,采用最简单的 range 分配策略。

分配逻辑是:把分区按数字排序,消费者也排序,然后把分区尽量均分。举例来说,6 个分区、3 个消费者,每个消费者分到 2 个分区。当有消费者离开或加入时,触发 rebalance,所有消费者停止消费,重新分配后再继续。

为了避免重复消费,作者让每个消费者记录自己消费到的 offset,提交到内部的一个“__consumer_offsets”主题(简化实现里只是一个普通日志文件)。rebalance 后,新接管分区的消费者从上次提交的 offset 继续读取。

这个过程暴露了 Kafka 消费者组的痛点:rebalance 期间会有短暂的停止服务,分配策略也可能导致分区分配不均。官方 Kafka 后来引入 sticky assignor 来减少不必要的分区迁移,自建版则停留在最原始的状态,让作者切实感受到生产环境需要多么精细的调优。

自建版本暴露的 Kafka 日志与索引设计权衡

日志存储是 Kafka 最核心的部分。自建时作者发现,如果只保留日志文件而不建索引,消费者想从特定 offset 读取就必须从头顺序扫描,效率极低。因此他增加了一个稀疏索引文件,每隔一定字节或一定消息数记录一次 offset 到文件位置的映射。

这种设计是官方 Kafka 的简化版。真实 Kafka 使用分段日志(每个日志文件达到一定大小就滚动),同时维护时间索引和 offset 索引,支持按时间戳查找。自建版只做了 offset 索引,省略了时间索引和日志清理策略。

日志清理是另一个重要权衡。作者没有实现基于时间的删除或 compaction,只是让日志无限增长。这让系统在长期运行后会消耗大量磁盘,但也避免了实现复杂清理逻辑带来的 bug。

通过这些取舍,作者看到 Kafka 的日志设计是一系列工程权衡的结果:顺序写带来高性能,稀疏索引平衡读速度与空间占用,分段日志方便清理。脱离实际场景很难理解为什么不采用数据库常见的 B-tree,而是坚持 append-only 加索引。

对中文开发者理解生产级消息系统的实际启发

这个自建项目对国内开发者最大的价值在于,它把官方文档里抽象的概念变成了可运行、可调试的代码。很多工程师读完 Kafka 源码还是觉得云里雾里,因为缺少一个能快速实验的玩具系统。

作者的轻量级实现不到一千行核心代码,却覆盖了日志追加、leader 选举、复制协议、消费者组协调等关键部分。中文开发者可以 fork 这个项目,在上面增加功能,比如实现 ISR、支持 Kafka 协议的 wire format,或者加上 metrics 监控。

更重要的是,这种从零构建的思路可以迁移到其他生产级系统:想懂 RocketMQ 的存储模型,就去写一个简化版 CommitLog;想懂 Pulsar 的 bookkeeper,就实现一个 quorum 复制层。每次动手,都能把“知道”和“理解”之间的差距大幅缩小。

在当前云原生和大数据岗位竞争激烈的环境下,真正理解底层机制的人更容易在面试中脱颖而出,也能在生产故障时更快定位问题。自建 Kafka 不是为了取代它,而是为了不再害怕它。

作者最后强调,这个项目不是教学玩具,而是个人学习路径上的里程碑。它证明了:最有效的学习方式,往往是把复杂系统拆到不能再拆,然后亲手拼起来。

(全文约 2150 字)

参考来源