2011 年,LinkedIn 的 Jay Kreps、Neha Narkhede、Jun Rao 在 NetDB 研讨会上发表了《Kafka: a Distributed Messaging System for Log Processing》。这篇论文里没有新算法,也没有新协议。它真正的主张只有一句:不要把消息系统当成队列来实现,把它当成一个只能往末尾追加的文件来实现。
之后十几年里,Kafka 几乎所有让人惊讶的特性——重放、多消费者、高吞吐、流批统一——都不是分别加上去的功能,而是这一个决定的推论。理解 Kafka,本质上就是理解"仅追加日志"这个抽象能推出多少东西,以及它推不出什么。
破除误解
误解一:Kafka 是一个"更快的 RabbitMQ"。
两者的数据结构根本不同。RabbitMQ 这类传统 broker 维护的是一个真正的队列:消息被投递、被确认(ack)之后从队列里删除,broker 需要为每一条消息记录"谁拿走了、有没有处理完"。Kafka 的 broker 不记录任何单条消息的消费状态。它只维护一个不断增长的文件,消费到哪里是消费者自己的事。所以 Kafka 天然没有单条消息确认、没有优先级队列、没有延迟投递、没有内建死信队列——不是"还没做",是这个数据结构里没有放这些东西的地方。反过来,RabbitMQ 也没法让你把上周的消息重新读一遍。
误解二:Kafka 保证消息全局有序。
Kafka 只保证单个分区内部的顺序。一个 topic 被切成若干分区,分区之间没有任何顺序关系。如果订单创建和订单支付两条消息落到不同分区,消费者完全可能先看到支付。要让相关消息有序,只能用同一个键(key)把它们哈希到同一分区——而这等于把这批消息的并行度锁死为 1。
误解三:开启"恰好一次"就不用写幂等逻辑了。
Kafka 的恰好一次语义(EOS)覆盖的是"从 Kafka 读 → 处理 → 写回 Kafka"这条闭环。一旦处理过程中要给外部系统写数据——发一封邮件、调一次支付接口、更新一张 MySQL 表——这个保证立刻失效,因为 Kafka 的事务管不到它管不着的系统。下面会详细拆开这一点。
磁盘上到底是什么
一个 topic 的一个分区,在 broker 的数据目录里就是一个文件夹。里面的东西朴素得有点意外:
/var/lib/kafka/data/orders-3/
00000000000000000000.log # 记录本体,只在末尾追加
00000000000000000000.index # 位移 → 文件字节位置(稀疏)
00000000000000000000.timeindex # 时间戳 → 位移(稀疏)
00000000000000368129.log # 下一个段,文件名就是段内首条记录的位移
leader-epoch-checkpoint
```.log 是日志段(segment)。写入只发生在最后一个"活跃段"的末尾;写满 log.segment.bytes(默认 1 GB)或超过 log.roll.hours(默认 168 小时)就滚动出新段。删除数据的方式也简单粗暴:整段删掉。默认保留期 log.retention.hours 是 168 小时,即 7 天。Kafka 没有"删除一条消息"这个操作——这不是疏漏,而是仅追加结构的必然:允许原地删改,顺序写的所有好处都没了。
位移(offset)是分区内单调递增的整数,从 0 开始,就是这条记录在这个分区里的序号。它不是 ID,不是时间戳,就是序号。.index 把位移映射到文件字节位置,但它是稀疏索引:默认每写满 index.interval.bytes(4096 字节)才记一条。查一个位移时,broker 二分索引找到最近的锚点,然后向后线性扫几 KB。用一点扫描换掉大部分索引体积——和 LSM 树里的稀疏块索引是同一个权衡。
一个日志,为什么同时是队列和流
这是整篇文章的枢纽。因为读取是移动一个指针而不是弹出一个元素,同一份数据可以被任意多方以任意速度、任意起点读取,互不干扰。于是:
| 想要的语义 | 在 Kafka 里怎么得到 |
|---|---|
| 工作队列(一条消息只被处理一次) | 一个消费者组,组内每个分区只分配给一个消费者 |
| 发布订阅(每个订阅方拿到全量) | 多个消费者组,各自维护各自的位移 |
| 重放 / 回溯重算 | 把位移重置到过去某点,重新消费 |
| 新系统冷启动灌数据 | 从位移 0 读到最新,然后自然转入实时 |
| 表(每个键的当前值) | 日志压缩:只保留每个键最后一次写入 |
消费者组是队列语义的来源:组内成员分摊分区,加一个消费者就多分几个分区,这就是"竞争消费"。多个组是发布订阅语义的来源:订单组、风控组、数仓组各读各的,谁也不影响谁。传统 broker 必须在实现层面把这两种模式做成两套东西,Kafka 只需要一套。
日志压缩(cleanup.policy=compact)则把日志变成表。压缩线程会在后台扫描旧段,对每个键只保留最后一条记录;写入 value 为 null 的"墓碑消息"表示删除。压缩后的 topic 读到最后,你手上就是一张完整的键值快照——这就是 Kafka Streams 里 KTable 的物理基础,也是所谓"流表二象性"的来源:流是表的变更序列,表是流的当前折叠。
分区:顺序保证的确切边界
分区是 Kafka 唯一的并行单位,也是它所有硬性限制的来源。
生产者默认用 murmur2(key) % 分区数 决定去哪个分区;键为空时,2.4 之后采用"粘性分区"策略,把一批消息尽量攒到同一个分区以提高批量效率。同一个键永远去同一个分区,因此同一个键的消息严格有序——这是 Kafka 提供的全部顺序保证。
几个容易踩的细节:
- 分区内有序也有前提。 生产者可以有多个请求同时在途(
max.in.flight.requests.per.connection,默认 5)。若关闭幂等生产者又开启重试,第一个请求失败重发就可能排到第二个后面,顺序被打乱。3.0 起enable.idempotence默认为true,broker 依据序列号识别乱序并纠正,才让"5 个在途 + 有序"同时成立。 - 分区数只能增不能减。 而增加分区会改变
key % 分区数的结果:扩容之后,同一个键可能落到新分区,而它的历史记录还留在旧分区。跨越扩容点的键级顺序保证就此断裂,这是生产环境最常见的隐性事故之一。 - 消费并行度的上限是分区数。 一个分区在一个组里只能被一个消费者持有。60 个分区的 topic,第 61 个消费者只能闲着。
- 分区多也有代价。 每个分区是若干个打开的文件句柄、一份副本同步状态、一份元数据;分区数上万后,故障切换时长和控制器压力会明显上升。
要全局有序?只能用一个分区,然后接受单机吞吐上限。Kafka 把"全序"这个昂贵的东西明码标价卖给你,而不是假装免费。
位移归消费者所有意味着什么
Broker 不追踪谁读到哪。消费者定期把自己的位移作为一条普通消息提交到内部 topic __consumer_offsets(默认 50 个分区,本身就是压缩 topic)。这条设计有几个直接后果:
Broker 是无状态的(就消费进度而言)。 它不需要为每个消费者维护游标、超时、重投计数,因此加消费者几乎不增加 broker 的负担。这是 Kafka 能同时服务大量消费方的根本原因。
交付语义由提交时机决定,不由 Kafka 决定。 先提交位移再处理,崩溃时消息丢失,是"至多一次";先处理再提交,崩溃时重复处理,是"至少一次"。默认 enable.auto.commit=true、每 5 秒自动提交一次——这个默认值给出的是一种模糊的至少一次,重复窗口最长 5 秒。想要确定的语义,必须关掉自动提交,自己控制提交点。
回溯是一等操作。 查看和修改进度就是两条命令:
# 看积压:CURRENT-OFFSET / LOG-END-OFFSET / LAG
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-processor --describe把这个组的消费位置拨回到某个时刻(需先停掉组内消费者) kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group order-processor --topic orders \ --reset-offsets --to-datetime 2026-08-01T00:00:00.000 --execute ```
代价是消费者落后没有反压。传统 broker 里队列堆积会挤爆内存,逼你正视问题;Kafka 里堆积只是 lag 数字变大,一切照常,直到消费者落后超过保留期——那时它要读的段已经被删了,客户端收到 OffsetOutOfRangeException,按 auto.offset.reset(默认 latest)跳到最新位置,中间的数据静默丢失。lag 监控在 Kafka 上不是锦上添花,是必需品。
为什么在磁盘上还能快
Kafka 把数据写在磁盘上却比很多内存队列还快,靠的是四件互相配合的事,没有一件是魔法。
一、顺序写。 Kafka 官方设计文档给出的对比是:在一组六块 7200 转 SATA 盘组成的 RAID-5 上,线性写约 600 MB/s,而随机写只有约 100 kB/s——相差六千倍以上。仅追加结构让 Kafka 永远只做前者。这个数字来自机械盘时代,SSD 上差距小得多,但顺序写依然更友好(对齐擦除块、减少写放大)。
二、依赖操作系统页缓存,而不是自己在 JVM 堆里做缓存。 Kafka 写入只是 write() 到页缓存,由内核决定何时刷盘;读取时如果数据还在页缓存里,根本不碰磁盘。好处是三重的:不占 JVM 堆因而没有 GC 压力;避免"应用缓存 + 页缓存"的双份拷贝;进程重启后缓存还在(页缓存属于内核)。典型的实时消费者读的都是刚写入的数据,命中率极高。
三、零拷贝。 把数据从文件发到网卡,常规路径要经过:磁盘 → 内核页缓存 → 用户态缓冲区 → 内核 socket 缓冲区 → 网卡,四次拷贝、多次用户态/内核态切换。Kafka 用 sendfile(Java 里是 FileChannel.transferTo)让数据直接从页缓存进 socket,跳过用户态。这能成立有个前提:broker 从不解析消息内容。消息对它就是一段不透明字节。一旦需要在服务端做过滤、转换、加密,零拷贝立刻失效——开启 TLS 后 Kafka 就用不上 sendfile,因为加密必须在用户态做。这是安全与吞吐之间一个真实存在、无法两全的取舍。
四、批量与端到端压缩。 生产者按 linger.ms 和 batch.size 攒批,整批压缩后写入;broker 以压缩后的原样存储;消费者拿到后自己解压。压缩只做一次、解压只做一次,中间的 broker 完全不介入。批量还把"每条消息一次系统调用"摊薄成"每批一次"。
这四件事有一个共同前提:broker 什么也不懂。 Kafka 的高性能不是优化出来的,是把功能拿掉换来的。
恰好一次:承诺了什么
0.11.0(2017 年 6 月,KIP-98)引入了两个机制。
幂等生产者:每个生产者启动时获得一个 PID(producer id)和 epoch,向每个分区发送的每一批消息带一个单调递增的序列号。Broker 为每个 (PID, 分区) 记住最近若干批的序列号,若收到的序列号不是"上一个 +1",就判定为重复并丢弃。这解决的是网络重试导致的重复写入——生产者发出请求、broker 写成功、ACK 在回程丢了、生产者重发,这条最常见的重复路径被封死了。
事务:生产者用一个稳定的 transactional.id 向事务协调器注册,可以把"写多个分区"和"提交消费位移"放进同一个原子单元。提交或中止时,Kafka 向涉及的每个分区写一条控制记录(事务标记)。设置 isolation.level=read_committed 的消费者会跳过被中止事务的消息,并且只能读到 LSO(Last Stable Offset)——即最早的未决事务之前的位置。
于是"读—处理—写"闭环成立:消费位移和处理结果一起提交,要么都生效,要么都不生效。Kafka Streams 把这套封装成 processing.guarantee 配置,2.6 起提供基于 KIP-447 的新实现(3.0 起名为 exactly_once_v2),旧实现随后被弃用移除。
没承诺什么
- 不承诺网络上只传一次。 重复的字节照样在网络上飞,只是被 broker 用序列号去重了。准确的说法是"恰好一次处理",不是"恰好一次投递"。当年围绕这个词的争论(有人搬出两将军问题反驳)根子就在这里。
- 不覆盖外部副作用。 事务只作用于 Kafka 内部的分区写入和位移提交。处理逻辑里调支付接口、发短信、写 Elasticsearch,这些动作 Kafka 既不知道也回滚不了。整条链路要恰好一次,下游必须自己幂等或自己支持事务,并与 Kafka 事务协同。
- 去重是有窗口的。 Broker 侧的 PID 状态不会永久保留(
producer.id.expiration.ms默认 1 天,transactional.id.expiration.ms默认 7 天)。生产者更换transactional.id、或状态过期后重新初始化,去重链条就断了。 - 不保护读端的错误重复消费。 如果消费者用
read_uncommitted(默认值),它照样能读到未提交甚至已中止的消息。 - 成本是实打实的。 每个事务在每个涉及的分区留下控制记录;
read_committed消费者被 LSO 卡住——一个卡死不提交的事务会让整条分区的下游停止推进,直到transaction.timeout.ms超时中止。这是生产上典型的"消费突然全停"事故成因。 - 不覆盖跨集群。 MirrorMaker 2 的跨集群复制是至少一次。
同一层面还有个常被忽略的配置陷阱:acks=all 只保证写入被当前 ISR(同步副本集合)中的所有副本确认。若 min.insync.replicas 仍是默认的 1,而副本一个个掉队后 ISR 只剩 leader,acks=all 就退化成 acks=1——leader 一挂数据就没了。三副本配置下应设为 2。
它在哪里会失效
- 不适合任务队列。 没有单条确认,一条处理不了的"毒丸消息"会卡住整个分区的队头,后面的消息全部阻塞。没有优先级、没有延迟投递、没有单条重试退避。这些场景 RabbitMQ、SQS 更合适。社区也承认这个缺口:4.0(2025 年 3 月)以早期访问形式引入了 KIP-932 共享组(share group),提供单条确认与重投,但它是补丁,不是原生形态。
- 不适合大量低吞吐 topic。 每个分区都有固定的文件与元数据开销,几万个稀疏 topic 的场景下资源浪费严重。
- 冷读会伤热路径。 一个从头重放的消费者会把大量旧数据拉进页缓存,挤掉实时消费者的热数据,让本来不碰盘的实时链路开始碰盘。分层存储(KIP-405,3.6 早期访问、3.9 生产可用)把冷数据卸到对象存储,正是冲这个问题去的。
- 存算耦合。 Broker 同时负责服务请求和持有分区数据,扩缩容意味着在节点间搬运 TB 级数据。这也是 Pulsar 用 BookKeeper 分离存储层的动机。
- 运维不是零成本。 元数据共识长期依赖 ZooKeeper,直到 KRaft 在 3.3(2022 年 10 月)标记生产可用、4.0 彻底移除 ZooKeeper 才收敛为单一进程模型。消费者组重平衡的抖动也是老问题,4.0 中 GA 的 KIP-848 新协议才系统性改善。
跨域连接
- 供应链:一个 topic 就是一条带库存的产线,保留期就是允许积压的成品仓。供应链研究的牛鞭效应说明,下游波动向上游传导时会被逐级放大;Kafka 的应对是切断这条传导路径——下游慢了不反压上游,而是让波动沉淀成磁盘上的库存。代价与实体库存完全同构:缓冲要占真实资源(7 天保留期就是 7 天的盘),而且库存会过期报废——消费者一旦落后超过保留期,损失的不是延迟而是数据本身。所以 lag 监控在这里等价于库龄监控,不是可选项。
- 工业工程与质量管理:消费者主动拉取而非 broker 推送,正是丰田看板的拉动式生产。推动式按上游节拍投料,下游跟不上就堆积甚至停线;拉动式让下游按自身实际能力取件,在制品数量由下游决定。Kafka 把这条原则直接写进了配置:
max.poll.records规定一次取多少件,max.poll.interval.ms(默认 5 分钟)规定超过多久没来取件就判定该工位失能、把分区重新分配给别人。这也解释了 Kafka 为什么没有优先级队列——拉动式产线上的工位不挑件,一旦允许挑件,先进先出的秩序就不复存在。 - 记忆系统:日志与物化视图的关系,对应心理学区分的情节记忆与语义记忆。情节记忆保存"何时发生了什么"的具体片段,语义记忆保存被抽提出来的稳定知识;后者由前者加工而成,却无法反推回前者。Kafka 的日志压缩做的正是这种不可逆抽提:只保留每个键的最后取值,把"这个账户余额变动过四十次"压成"这个账户余额是 X"。人类记忆的教训在这里逐字适用——一旦只剩语义层,任何需要复盘中间过程的审计、排障与重算都做不成了。所以压缩 topic 与保留完整历史的 topic 必须按用途分开设计,而不是当成同一个东西的两档参数。
- 什么是真理:把日志称作"单一真相源",在认识论上是一次明确的站队。符合论认为一个陈述之为真在于它符合外部事实,融贯论认为真在于它与整个信念系统相容。事件溯源架构选的是后者:数据库里那个余额之所以"对",不是因为它符合某本外部账本,而是因为它与日志重放的结果融贯。这个选择有非常具体的成本——日志之上不再有更高的裁判,写错的事件不能改写,只能追加一条补偿事件去抵消;于是系统里永远同时存在"曾经的真"与"现在的真",任何读取都必须声明自己要的是哪一个。
- 热力学定律:仅追加是把不可逆性直接写进了数据结构。热力学第二定律刻画的是宏观过程的方向性:由完整历史可以推出当前状态,由当前状态推不回历史。Kafka 的日志正是这种不对称的工程化——写入只发生在末端,读取只是移动指针,任何"撤销"都只能表现为一次新的追加。顺序写之所以比随机写快出几个数量级(官方设计文档给出的对比是同一组机械盘上约 600 MB/s 对约 100 kB/s),根源同样在于放弃了原地折返的自由:速度不是优化出来的,是拿可逆性换来的。
参考文献
- Kreps, J., Narkhede, N., Rao, J. Kafka: a Distributed Messaging System for Log Processing. NetDB Workshop, 2011.(原始论文,日志抽象的一手出处)
- Kreps, J. The Log: What every software engineer should know about real-time data's unifying abstraction. LinkedIn Engineering, 2013.(把"日志即统一抽象"讲透的长文)
- Apache Software Foundation. Apache Kafka Documentation — Design.(顺序写与随机写的 600 MB/s 对 100 kB/s 对比、页缓存与 sendfile 论证的官方出处)
- Apache Software Foundation. KIP-98: Exactly Once Delivery and Transactional Messaging, 2017;KIP-447: Producer Scalability for Exactly Once Semantics, 2020;KIP-932: Queues for Kafka, 2025.(幂等、事务与共享组的一手提案)
- Apache Software Foundation. Apache Kafka 4.0.0 Release Announcement, 2025 年 3 月 18 日.(KRaft 唯一模式、KIP-848 GA、KIP-932 早期访问)
- Kleppmann, M. Designing Data-Intensive Applications, Ch. 11 "Stream Processing". O'Reilly, 2017.(日志、流表二象性与交付语义的系统梳理)
延伸阅读
- Narkhede, N., Shapira, G., Palino, T. Kafka: The Definitive Guide. O'Reilly(第 2 版 2021)——分区、位移、事务的工程细节。
- Confluent. Exactly-once Semantics Are Possible: Here's How Apache Kafka Does It. 2017——EOS 实现者视角的解释与争论回应。
- Kleppmann, M. Turning the Database Inside-Out with Apache Samza. Strange Loop, 2014——把数据库内部的日志翻到架构中心的思想来源。