Chapter 02
原理与设计取舍
上一章把 Kafka 框成一条按分区切分、可重放的提交日志,给了一套能用来描述系统的词汇——topic / partition / offset / leader / follower / consumer group。这章把那条日志拆开:它在磁盘上怎么存、在副本之间怎么复制、被消费组怎么读、跨重启与失败怎么保证语义。每个机制都配一张备选方案对比表——记住"为什么不选另一条路",比记住"选了什么"更经得起面试官追问。
本章你将建立的 schema
- 日志为什么快:append-only + segment 索引 + page cache + 零拷贝,以及它们各自的失效边界
- 顺序与并行同一个单位:partition——key 哈希定分区、sticky 批处理、并行度上限
- 持久性的真实来源:LEO / HW / ISR 三者关系,以及"
acks=all不等于所有副本" - 再平衡的三代协议:eager → cooperative → KIP-848,以及"心跳活着也会被判死"
- 投递语义的真实边界:幂等 producer、事务、EOS 只在 Kafka 闭环内、只在同会话同分区成立
- 元数据也是一条日志:KRaft 用 Raft 复制
__cluster_metadata,取代 ZooKeeper
这张图是本章的地图。后面六节各放大其中一块:§2.1 存储看 leader 那条日志在磁盘上长什么样;§2.2 分区与顺序看一条 topic 为什么要切成多个 partition;§2.3 副本与 ISR看 follower 的 fetch 如何决定"哪些数据对 consumer 可见";§2.4 消费组与再平衡看右边那个 consumer group 内部怎么分配分区;§2.5 投递语义看一条记录从写到读"恰好一次"到底意味着什么;§2.6 KRaft看底部那条元数据日志。
2.1存储原理:日志为什么快
每个 partition 是一个 append-only 文件序列,切成定长 segment、配稀疏索引,靠操作系统 page cache 和零拷贝把磁盘的顺序 I/O 跑到接近内存的速度。
运行方式
一个 partition 在磁盘上不是一个大文件,而是一串segment(日志段,由 log.segment.bytes 控制大小,默认 1GB)。任何时刻只有最后一个 active segment 在被追加写;写满就滚动出新的一个,旧的变成只读。每个 segment 旁边配两个索引文件:.index 把 offset 映射到该 segment 内的字节位置,.timeindex 把时间戳映射到 offset。两个索引都是稀疏的(不是每条记录都建索引,而是每隔几 KB 建一条),所以查一条记录是先二分索引定位到附近、再顺序扫一小段——O(log N) 定位而不是全文件扫描。
关键在于读写都是顺序 I/O:写永远是追加到文件末尾,读是从某个 offset 起连续往后。Kafka 不在 JVM 进程内缓存消息,而是直接依赖操作系统的 page cache——写入的数据先进 page cache(由 OS 异步刷盘),消费者读最近的数据时几乎总是直接命中 page cache,根本不碰磁盘。消费者读取走零拷贝(sendfile / FileChannel.transferTo()):数据从 page cache 直接进网卡缓冲区,不经过 JVM 用户态——省掉了 "内核态→用户态→内核态" 的两次拷贝和上下文切换。
# topic "orders" 的分区 0 在 broker 数据目录下:
$ ls -lh /var/kafka-logs/orders-0/
00000000000000000000.log # segment:实际记录(base offset = 0)
00000000000000000000.index # 稀疏 offset→字节位置 索引
00000000000000000000.timeindex # 稀疏 时间戳→offset 索引
00000000000000368912.log # 下一个 segment,base offset = 368912
00000000000000368912.index
00000000000000368912.timeindex
leader-epoch-checkpoint # leader 任期记录(截断时用,见 §2.3)
# 文件名 = 该 segment 第一条记录的 offset(base offset)。
# 查 offset 400000:文件名二分 → 落在 368912.log → .index 二分定位到附近字节 → 顺序扫到 400000。
保留策略有两种。默认是按时间/大小删整段(retention.ms / retention.bytes,到期把最旧的 segment 整个删掉)。另一种是 log compaction(日志压缩):对每个 key 只保留最新一条 value,老版本被清理——把 topic 变成一份"可恢复的状态快照"。写一条 value 为 null 的记录叫 tombstone(墓碑),表示"这个 key 被删除",compaction 在保留窗口后把该 key 连同墓碑一起清掉。__consumer_offsets(消费位移)和 Kafka Streams 的 changelog 就是靠 compaction 实现的——这也是第 1 章"日志可重放"能延伸成"日志可作状态存储"的机制。
备选方案对比
| 方案 | 优势 | 为什么没选 |
|---|---|---|
| JVM 堆内缓存(进程内对象池) | 访问对象零序列化、命中即返回 | 与 OS page cache 双份缓冲(同一份数据存两遍);几十 GB 缓存进堆会制造灾难性 GC 停顿;进程一重启缓存全冷 |
| 堆外 + 自管缓存(off-heap 自己淘汰) | 绕开 GC、可控大小 | 等于重写一个比 OS 还差的 page cache;故障/重启仍丢缓存;复杂度全压给应用 |
| OS page cache(不在进程内缓存) | 无 GC 压力、32GB 机器约 28-30GB 可用作缓存、进程重启缓存仍在内核、零拷贝可直接从它发网卡 | 选中 |
| 方案 | 优势 | 为什么没选 |
|---|---|---|
| B-Tree / LSM(数据库式索引) | 支持任意 key 随机读写、范围查询 | 写要随机寻道 + 页分裂 + 加锁,O(log N) 且放大写;Kafka 根本不需要"按任意 key 改某条记录",只需要"按 offset 顺序追加和顺序读" |
| 每条消息一个文件 | 删除单条简单 | 海量小文件耗尽 inode 和文件句柄;丢失顺序 I/O 的全部好处 |
| append-only 日志 + 稀疏索引 + 定长 segment | 追加 O(1)、顺序磁盘 I/O 接近内存带宽、删除按段 O(1)、稀疏索引省内存 | 选中 |
顺序磁盘 I/O 的吞吐接近内存、远高于随机 I/O(机械盘上差几个数量级,SSD 上也有明显差距)。进程内缓存会与 page cache 重复存储、还加重 GC,所以把缓存这件事整个交给 OS。append-only 追加是 O(1),对比 B-Tree 的 O(log N) 随机寻道加锁。compaction 让一个 topic 能当 changelog / 状态快照用,而不只是消息管道。
带来的代价 / 失效模式
① TLS 绕过零拷贝。sendfile 要求内核能直接把文件字节送上网卡。一旦开启传输加密(broker 间或 client-broker TLS),数据必须先进用户态加密再发出——零拷贝失效,CPU 和延迟都上升。这是"开了 TLS 吞吐就掉一截"的根因。
② 冷读污染 page cache。page cache 的前提是"消费者读的就是刚写的热数据"。一个严重滞后的消费者去读几小时前的旧数据,会把磁盘上的冷 segment 拉进 page cache,挤掉正在服务热生产者的页——一个掉队的消费者能拖慢整个 broker。Tiered Storage(KIP-405,3.9 GA)把冷数据下沉到对象存储正是为缓解这点。
③ compaction 是最终去重,不是即时。compaction 只在脏数据比例超过 min.cleanable.dirty.ratio(默认 0.5)时才触发;tombstone 在 delete.retention.ms 内仍然存在以便下游消费者看到删除。所以"写了新值老值就没了""写了墓碑 key 立刻消失"都是错的——是最终一致的去重。
你给一个 topic 开了 TLS、又有一个消费者从 offset 0 全量重刷历史数据。这两件事各自怎么影响 broker 的吞吐?它们打击的是同一个机制吗?
展开答案(先停 10 秒再点)
打击的是两个不同机制。TLS 关掉的是零拷贝:数据被迫走用户态加密,CPU 上升、每条消息多两次拷贝。全量重刷打击的是 page cache 命中率:从 offset 0 读的全是冷 segment,把热生产者的页挤出缓存,命中率暴跌、磁盘随机读上升。
设计洞察:Kafka 的"快"建立在两个独立假设上——零拷贝(数据不进用户态)+ 热数据局部性(读的就是刚写的)。TLS 破坏前者,滞后消费者破坏后者。理解这两条,就能解释绝大多数"Kafka 突然变慢"的现场。
2.2分区与顺序:并行和有序的同一个单位
partition 同时是顺序的单位和并行的单位——同一个 key 永远落同一个 partition(murmur2 哈希)因而有序,不同 partition 之间并行;顺序只在分区内成立。
运行方式
一条记录写到哪个 partition,由 producer 端决定。带 key 时:partition = murmur2(key) % 分区数——相同 key 永远算出同一个 partition,因此同一个 key 的记录在该分区内严格有序(同一订单的所有事件按写入顺序排好)。不带 key 时,3.0 以后默认用 sticky partitioner(KIP-480):先把一个分区的 batch 填满再换下一个分区,而不是旧的 round-robin 逐条轮询。批越满,每批的固定开销摊得越薄,端到端延迟约减半。
这就是第 1 章里"顺序只在分区内保证"的机制根源:offset 是每个分区独立递增的序号,跨分区没有全局序。消费侧的并行度也由分区数决定——一个 consumer group 里,一个 partition 最多被一个消费者持有,所以消费者数超过分区数就有人闲置(详见 §2.4)。
| 方案 | 优势 | 为什么没选 |
|---|---|---|
| 全 topic 全局总序(所有消息一条线排序) | 消费者看到的就是绝对时间序,推理最简单 | 全局有序需要单一写入点序列化所有写——无法水平扩展,吞吐被一台机器锁死;一个分区杀掉整个并行 |
| 完全不保证顺序(纯负载均衡分发) | 写入和消费都可无限并行 | 无法表达"同一实体的事件有先后"——订单"创建→支付→发货"可能乱序到达,业务无法处理 |
| 每分区有序(按 key 哈希分区) | 顺序廉价(分区内天然有序)、并行度可扩到分区数、同 key 同序满足绝大多数业务 | 选中 |
| 方案 | 优势 | 为什么没选 |
|---|---|---|
| round-robin(逐条轮询分区) | 分区间负载绝对均匀 | 每条记录可能进不同分区的 batch,导致大量半满小批发送——请求数多、吞吐低、延迟高 |
| 随机分区(每条随机选) | 实现简单、长期均匀 | 同样打散 batch,且短期可能不均;没有解决批处理问题 |
| sticky partitioner(KIP-480,先填满再换) | batch 更满 → 请求更少 → 延迟约减半;长期看分区仍大致均匀 | 选中 |
全局有序与水平扩展是矛盾的:要么有单写入点(不能扩展),要么放弃全局序。Kafka 选择把"有序"的粒度降到分区——业务真正需要的几乎都是"同一实体内有序"(同一用户、同一订单),用 key 哈希就能廉价拿到,同时分区数给了并行的旋钮。sticky 则是在"无序消息"这个子问题上,用牺牲瞬时均匀换批处理效率。
带来的代价 / 失效模式
① 并行度上限 = 分区数。一个 consumer group 的消费并行度卡在分区数。分区开少了,加再多消费者也无法提速(多出来的消费者空转);分区开多了,抬高 controller 元数据负担、故障切换时间、文件句柄消耗(每 segment 占 2 个文件,默认 ulimit 1024 很容易触顶)。
② 加分区会永久打乱 key 顺序。分区数只能加不能减。而一旦给带 key 的 topic 加分区,murmur2(key) % 新分区数 的结果变了——同一个 key 的新记录可能落到与历史记录不同的分区,跨分区的全局顺序无从保证,所有依赖按 key 顺序的下游消费者被静默破坏。要扩容只能保守预估,或新建 topic 迁移。
③ 热 key 倾斜压垮单分区。哈希只保证 key 均匀,不保证流量均匀。若 90% 流量集中在少数 key(如按租户分区而某个大租户占大头),这些 key 全挤进同一个分区,该分区的 leader 和它的消费者被打爆,其余分区闲置。Kafka 不会自动均衡倾斜——只能选高基数的 key(用 orderId 而不是 userId)。
一个有 6 个分区、按 userId 分区的 topic,运行半年后为了扩容把分区加到 12 个。下游有个消费者依赖"同一用户的操作严格有序"。它会立刻报错吗?如果不报错,问题什么时候暴露?
展开答案(先停 10 秒再点)
不会立刻报错——这正是它危险的地方。加分区是个成功的管理操作,没有任何异常。但从加分区那一刻起,某个 userId 的新记录算出的分区(% 12)可能不同于它历史记录所在的分区(% 6)。于是同一用户的事件流被劈到两个分区,两个分区由不同消费者并行消费,先后顺序彻底丢失。
问题在"恰好某个用户的新旧事件被并发处理、且顺序敏感"时才暴露——可能是几天后一次诡异的状态错乱,而且极难复现。所以 ctx 里把它列为"静默破坏":没有报错的破坏最贵。生产上对有 key 且顺序敏感的 topic,要么一开始就把分区数开够,要么新建 topic 做带迁移的切换。
2.3副本与 ISR:持久性的真实来源
每个分区一个 leader 加若干 follower,follower 像消费者一样主动 fetch;HW(高水位)= ISR 中最小的 LEO,消费者只能读到 HW——已 ack 的记录未必立刻可见,而持久性来自 min.insync.replicas 不是 acks。
运行方式
每个分区有一个 leader 副本和若干 follower 副本。所有读写都走 leader;follower 主动向 leader fetch(和普通消费者用同一套 fetch 机制,只是 fetch 的是分区日志去复制)。两个关键水位:
- LEO(Log End Offset):某个副本下一条要写入的 offset,也就是它当前日志的末尾。leader 的 LEO 总是最靠前,follower 的 LEO 落在后面(差着 fetch 的延迟)。
- HW(High Watermark,高水位):ISR 集合里所有副本 LEO 的最小值,即"已经被所有同步副本完整复制到的最高 offset"。
消费者只能读到 ≤ HW 的记录。HW 之上、leader 已写但还没被所有 ISR 复制完的那段,对消费者不可见——因为那段一旦 leader 故障、由某个 follower 顶上,可能会消失。HW 这道门控保证消费者读到的都是"即使现在 leader 挂了也不会丢"的记录。
ISR(In-Sync Replicas,同步副本集)是当前跟得上 leader 的副本集合(含 leader 本身)。某个 follower 若超过 replica.lag.time.max.ms(默认 30s)没追上 leader 的末尾,就被踢出 ISR;追上后再加回来。ISR 是动态可变的——这是理解后面所有失效模式的关键。
acks 与 min.insync.replicas 的真实关系
producer 的 acks 决定"写入要等到什么程度才算成功":
acks=0:发出去就算成功,不等任何确认——最快,leader 没写成也不知道,丢数据风险最高。acks=1:leader 写入本地日志就 ack——leader 写完、还没被任何 follower 复制就宕机,那些已 ack 的记录永久丢失。acks=all(即acks=-1):等到当前 ISR 里所有副本都复制完才 ack。
这里是面试最容易露馅的一点:acks=all 等的是"所有 ISR 成员",不是"所有副本"。ISR 是会缩的——如果两个 follower 都掉队被踢出 ISR,ISR 只剩 leader 一个,此时 acks=all 实际只写到 leader 一台,悄悄退化成了 acks=1,下一次故障就丢。
真正的持久性闸门是 min.insync.replicas:它要求 ISR 至少有这么多成员在线,否则 producer 的写入直接被拒(收到 NotEnoughReplicas / NotEnoughReplicasAfterAppend)。所以持久性来自 min.insync.replicas,不是 acks——acks=all 负责"等所有 ISR",min.insync.replicas 负责"ISR 不许缩到不安全"。两者必须配合:典型生产配置是 replication.factor=3 + min.insync.replicas=2 + acks=all,含义是"至少 2 个副本拿到才算写成功,否则宁可拒写"。这套组合留到 §4.1 数据丢失里逐条对照。
| 方案 | 优势 | 为什么没选(或选中理由) |
|---|---|---|
| 多数派 quorum(Raft/Paxos 式,2f+1 副本容忍 f 故障) | 提交只需多数派,单个慢副本不拖累;成员是否在线靠投票自然处理 | 要容忍 f 个故障需 2f+1 个副本(容忍 2 个故障要 5 副本)——存储和带宽成本翻倍多;对"写多读多"的日志太贵 |
| 全副本同步(写必须等全部副本) | 任意 f<RF 故障都不丢,逻辑最简单 | 最慢的那个副本决定写延迟;任何一个副本卡住或宕机就停写,可用性差 |
| ISR 可变法定集(f+1 副本 + 动态 ISR) | 容忍 f 故障只需 f+1 副本(容忍 2 个故障只要 3 副本);掉队副本被移出 ISR 不拖慢写;用 min.insync.replicas 显式调持久性/可用性平衡 |
选中(更省副本,代价见下) |
多数派 quorum 为了"无需等慢副本"付出了 2f+1 的副本成本。Kafka 的洞察是:日志数据量大、副本成本敏感,而它已经有了"踢掉慢副本"的机制(ISR)——所以用 f+1 个副本 + 动态 ISR 就能容忍 f 个故障,省掉将近一半副本。代价是把"持久性 vs 可用性"的旋钮(min.insync.replicas、unclean.leader.election)交给运维显式决定,而不是协议帮你定死。
带来的代价 / 失效模式
① HW 滞后一个 fetch 轮。HW 的推进依赖 follower 下一轮 fetch 上报自己的 LEO。所以一条记录即使已经被所有 ISR 复制、已经 ack 给 producer,consumer 也要再等一个 fetch RTT 才能看到它(如图 2.1)。"已 ack 未必立刻可见"是设计使然,不是 bug。
② acks=all + min.insync.replicas=2 在 RF=3 下:挂 1 台仍可写,挂 2 台分区停写。这是故意用可用性换持久性——剩 1 个副本时与其欠复制地写(将来丢),不如拒写让 producer 知道。能不能接受停写,是个业务决策。
③ unclean leader election 静默丢数据。若开启 unclean.leader.election.enable=true,当 ISR 里的副本全挂、只剩一个落后的非 ISR 副本时,Kafka 允许它当 leader 以恢复可用性。但它的日志比之前已提交的短——比如 producer 已 ack 到 offset 100,这个落后副本只到 80,于是 81-100 被静默截断、永久丢失。这就是 leader-epoch-checkpoint 文件存在的场景:靠 leader 任期来判断哪段日志需要截断。默认值在新版已是 false(宁可分区不可用也不丢)。
某 topic 配置 replication.factor=3、acks=all,但没有设 min.insync.replicas(用默认值 1)。平时一切正常。某天两个 follower 因网络抖动同时掉出 ISR,运维没注意。这段时间 producer 的写入安全吗?
展开答案(先停 10 秒再点)
不安全,而且是看不见的不安全。min.insync.replicas=1 意味着 ISR 只剩 leader 一个时仍允许写。两个 follower 掉出 ISR 后,ISR={leader},此时 acks=all 的"all"就是 leader 自己——每次写实际只落到一台,等价于 acks=1。producer 收到的全是成功 ack,毫无察觉。只要这期间 leader 宕机,这批"成功"的记录全部丢失。
洞察:acks=all 单独配置是不够的——它的强度被 ISR 的当前大小绑架。min.insync.replicas=2 才能在 ISR 缩到 1 时让写入被拒而不是悄悄降级。这正是 §2.3 的核心命题"持久性来自 min.insync.replicas"的实战形态,也是 04 章数据丢失类的头号案例。
2.4消费组与再平衡:分区如何分配,成员如何被判死
broker 端的 group coordinator 通过 JoinGroup/SyncGroup 把分区分给组内消费者;再平衡协议从 eager(全停)演进到 cooperative(只动该动的)再到 KIP-848(broker 端驱动),而成员存活由 poll() 证明、不是心跳。
运行方式
每个 consumer group 由 broker 端一个 group coordinator 管理。新成员加入或老成员离开时触发再平衡,分两步:JoinGroup(所有成员向 coordinator 报到,coordinator 选出其中一个当 group leader)→ SyncGroup(group leader 在客户端算出"谁拿哪些分区"的分配方案,回传给 coordinator 下发)。注意分配逻辑跑在客户端的 group leader 上,broker 不关心用什么策略分——这让分配策略可插拔。
消费位移(offset)提交到内部 topic __consumer_offsets(一个 compaction 的 topic,每个 (group, topic, partition) 只保留最新 offset)。提交可以自动(enable.auto.commit,按 auto.commit.interval.ms 定时)或手动(commitSync() / commitAsync())。
成员存活判定有两条独立的线,这是最容易混淆的地方:
- 后台心跳线程按
heartbeat.interval.ms给 coordinator 发心跳,session.timeout.ms内没心跳就判死——这检测的是"进程/网络是否还活着"。 - 但真正证明"消费者还在干活"的是调用
poll()。两次poll()间隔超过max.poll.interval.ms(默认 5 分钟),coordinator 判定这个成员卡死,把它踢出组触发再平衡——即使它的心跳线程还在后台正常跳。
三代再平衡协议
eager(急切式):每次再平衡,所有成员先放弃全部分区(一个 stop-the-world 屏障),等新分配下来再重新认领。简单、分配干净,但再平衡期间整个组停止消费。
cooperative / incremental(协作增量式,2.4+):遵循"没必要动的资源不要停"——只收回需要转移的那部分分区,没变动的分区继续消费;分两轮收敛。
KIP-848(4.0 GA):把再平衡逻辑从客户端移到 broker 端,由 coordinator 增量驱动,不再有客户端侧的全局屏障,加入/退出更平滑。经典 eager 协议计划在 5.0 移除。
static membership(静态成员,KIP-345):给消费者配一个固定的 group.instance.id。配了之后,成员短暂重启(如滚动发布)能认领回原来的分区而不触发再平衡——把"重启"和"成员变更"解耦。
| 方案 | 优势 | 为什么被取代 / 选中 |
|---|---|---|
| eager(急切式,stop-the-world) | 协议简单、分配干净,一轮收敛 | 每次再平衡整组暂停消费;成员越多、分区越多停顿越久(900 任务的组启动要 12-14 分钟)——5.0 将移除 |
| cooperative / incremental(协作增量式,2.4+) | 只收回需迁移的分区,其余继续消费;同样规模启动从十几分钟降到约 60 秒 | 分两轮收敛、协议更复杂,且仍在客户端驱动——被 KIP-848 进一步改进 |
| KIP-848(4.0 GA,broker 端驱动增量) | 再平衡逻辑移到 broker 端、增量推进、无客户端全局屏障;加入/退出最平滑 | 选中(当前方向;旧 eager 5.0 移除) |
eager 的 stop-the-world 屏障保证了分配的简单和正确,但代价随集群规模线性放大——大组的每次扩缩容都是一次全组停顿。cooperative 把原则换成"不动的资源不要停",用协议复杂度换可用性。KIP-848 再进一步:客户端驱动的协商在大规模下本身就脆弱(客户端版本不一、网络分区),把它收归 broker 端能更可控地增量推进。三代的主线是同一个:把再平衡的"暂停面"越缩越小。
带来的代价 / 失效模式
心跳活着,poll 超时仍被判死 → 重复处理。这是最反直觉的失效。后台心跳线程和业务处理是两回事:消费者在 poll() 之后花了 6 分钟处理这批记录(一个慢的 DB 批写、一次卡住的 HTTP 调用),心跳一直正常跳,但因为超过了 max.poll.interval.ms(默认 5 分钟)没有再调 poll(),coordinator 已经把它判死、把它的分区分给了别人。等它处理完想提交 offset,发现自己已被踢出组,提交失败;那批记录被新 owner 重复处理了一遍。
这是再平衡风暴的种子:处理慢 → 被踢 → 重新加入 → 又处理慢 → 再被踢,吞吐可以归零。缓解:调小 max.poll.records(每次少拿点)、把重活卸到工作线程并用 pause()/resume()、开启协作式再平衡 + static membership。
有人把 session.timeout.ms 调到 5 分钟,理由是"防止消费者因为偶尔的网络抖动被误判死"。处理逻辑每批要跑 8 分钟。会发生什么?调 session.timeout.ms 解决得了吗?
展开答案(先停 10 秒再点)
解决不了,因为他调错了旋钮。session.timeout.ms 管的是心跳那条线——进程/网络是否存活。但 8 分钟的处理超时触发的是另一条线:max.poll.interval.ms(默认 5 分钟)。处理跑 8 分钟没调 poll(),不管心跳多正常,到 5 分钟就被判定为"卡死"踢出组。
正确做法是认清两条线各管什么:要容忍长处理,调大的是 max.poll.interval.ms(或减小 max.poll.records 让每批处理更快返回 poll()),而不是 session.timeout.ms。把心跳超时调到 5 分钟反而有害——真宕机的消费者要拖 5 分钟才被发现。理解"存活由 poll() 证明、不是心跳"是这道题的全部。
2.5投递语义:幂等、事务与 EOS 的真实边界
at-least-once 是默认;幂等 producer 用 (PID, 分区) + 序列号在 broker 端去重并强制保序;事务用 transactional.id + epoch 把"输出写 + offset 提交"做成原子——但 exactly-once 只在 Kafka 闭环内、且幂等只在同一会话同一分区成立。
三种语义的来源
- at-most-once(至多一次):先提交 offset 再处理。处理崩了 offset 已提交,这条不会重投——但也可能丢。
- at-least-once(至少一次,默认):先处理再提交 offset。处理完崩在提交前,重启会重投这条——不丢,但可能重复。
- exactly-once(精确一次):靠幂等 producer + 事务达成,且只在特定边界内成立(见下)。
幂等 producer:PID + 序列号
开启 enable.idempotence=true(4.0 起默认开)后,每个 producer 会话拿到一个 PID(Producer ID)。producer 给发出的每条记录按 (PID, 分区) 维度打一个单调递增的序列号。broker 为每个 (PID, 分区) 记住"已接受的最高序列号",只接受序列号恰好 +1 的 append:
- 收到重复的序列号(重试导致的重发)→ 直接丢弃,去重。这解决"broker 写成功但 ack 在网络上丢了、producer 重发"导致的重复。
- 收到有空洞的序列号(比预期跳号)→ 抛
OutOfOrderSequenceException。
因为只接受恰好 +1,幂等 producer 顺便强制了顺序:它会拒绝乱序到达的 batch,所以同时解决了"max.in.flight.requests > 1 时重试导致消息乱序"这个老问题(幂等下 in-flight 可达 5 仍保序)。
事务:transactional.id + epoch + 两阶段提交
幂等只覆盖单个 producer 会话、且只去重不跨"多分区写 + offset 提交"。事务(配 transactional.id)补上这两块:
- 跨重启的身份与僵尸隔离:
transactional.id是稳定的(重启不变),每次初始化会把 epoch 加一。旧实例(僵尸)带着旧 epoch 再来写会被拒——避免"以为死了其实没死"的旧进程污染数据。 - 原子的"输出 + offset":transaction coordinator 跑两阶段提交,给事务涉及的每个分区写一个 commit / abort marker。"消费-转换-生产"里,把下游输出和上游 offset 提交放进同一个事务——要么都成功,要么都回滚,不会出现"输出写了但 offset 没提交"导致的重复。
- 读侧配合:消费者设
isolation.level=read_committed后,会跳过 aborted 和未决事务的数据,且读取不能越过 LSO(Last Stable Offset,最后稳定位移)——即最早的未决事务之前的位置。
// "消费 → 转换 → 生产" 的 EOS 闭环:输出写 + 上游 offset 提交在一个事务里
Properties p = new Properties();
p.put("transactional.id", "order-enricher-1"); // 稳定 ID,重启不变 → epoch 隔离僵尸
p.put("enable.idempotence", "true"); // 事务隐含开启幂等
KafkaProducer<String, String> producer = new KafkaProducer<>(p);
producer.initTransactions(); // 领 PID、把 epoch +1,挤掉僵尸旧实例
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(200));
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> r : records) {
producer.send(new ProducerRecord<>("orders-enriched", r.key(), enrich(r.value())));
}
// 关键:上游 offset 也写进同一个事务,而不是 consumer.commitSync()
producer.sendOffsetsToTransaction(offsetsOf(records), consumer.groupMetadata());
producer.commitTransaction(); // 写 commit marker,原子生效
} catch (KafkaException e) {
producer.abortTransaction(); // 写 abort marker,read_committed 端跳过
}
}
| 方案 | 优势 | 为什么没选 / 选中 |
|---|---|---|
| at-least-once + 消费侧幂等(下游用主键去重 / upsert) | producer/consumer 配置简单、无事务开销;对外部系统(DB)天然适用 | 把去重责任推给每个下游,每个消费者都要自己实现且实现正确;不解决"Kafka 到 Kafka 的多分区原子写" |
| at-most-once(先提交后处理) | 绝不重复、实现最简单 | 用丢失换"不重复"——崩在处理前那条就没了,绝大多数业务不可接受 |
| 事务 EOS(幂等 producer + transactional.id + read_committed) | Kafka 闭环内"消费-转换-生产"原子且精确一次;跨会话靠 epoch 存活;服务端去重让客户端简单 | 选中(边界与代价见下) |
① 幂等只在"同一个 (PID, 分区) 且同一 producer 会话"内去重。producer 一旦崩溃重启,会领到一个新 PID,broker 不认识它和旧 PID 的关系——跨重启不去重。要跨会话精确一次,必须用事务:transactional.id 靠 epoch 跨重启保持身份。所以"开了 enable.idempotence=true 就端到端精确一次"是被普遍过度宣称的误解。
② "exactly-once" 只在 Kafka 的"消费-转换-生产"闭环内成立,不覆盖外部副作用。事务能保证的是"Kafka 输出 + Kafka offset 提交"这件事原子。但事务里如果还写了数据库、调了 REST、发了邮件——这些外部动作不在事务范围内,事务回滚不会撤销它们,重试会再触发一次。要把外部系统也纳入精确一次,得在外部侧做幂等(如带幂等键的 upsert),Kafka 事务本身管不了。
把去重和隔离做在服务端(broker 维护 PID→seq、coordinator 管 marker),客户端就能保持简单——不用每个应用都重写一遍去重逻辑。而事务复用了已有的副本机制(事务日志、marker 都是普通的、被复制的日志记录)来保证自身的持久性,没有另起炉灶。代价是 marker + LSO 带来的延迟与队头阻塞(见下)。
带来的代价 / 失效模式
① marker + LSO 增加延迟。每个事务要额外写 commit/abort marker,read_committed 消费者要等事务定下来才能越过 LSO 读取——实测约增加 3%(@100ms 提交间隔)的延迟。提交越频繁,marker 开销占比越高。
② 长事务造成队头阻塞。read_committed 消费者读取不能越过 LSO。一个迟迟不提交的长事务会把 LSO 钉在原地,卡住该分区所有 read_committed 读者——哪怕后面已经有大量已提交的数据,也读不到。事务要短。
一个服务"从 topic A 消费 → 写一行到 MySQL → 往 topic B 生产",全程用了 Kafka 事务(sendOffsetsToTransaction + commitTransaction)。开发者声称"现在端到端精确一次了"。MySQL 那行会不会重复?
展开答案(先停 10 秒再点)
会重复。Kafka 事务保证的是"topic B 的输出 + topic A 的 offset 提交"这两件 Kafka 内部的事原子。MySQL 的写不在事务范围内——它是个外部副作用。设想:写完 MySQL、还没 commitTransaction 时进程崩溃。Kafka 这边事务未提交、会回滚(topic A 的 offset 没推进),重启后这批记录被重新消费,于是 MySQL 那行被写第二遍。Kafka 的回滚撤不掉已经发生的 MySQL 写。
洞察:EOS 的边界是"Kafka 闭环"。一旦事务里掺进任何外部系统,就要在那个系统侧另做幂等(MySQL 用幂等键 upsert / 唯一约束)。"用了 Kafka 事务 = 全链路精确一次"是 06 章会专门考的红旗答案。
2.6KRaft:元数据本身也是一条日志
集群元数据被建模成一条内部日志 __cluster_metadata,由 controller quorum 用 Raft 复制;active controller 是唯一写入者,其余节点回放日志并把元数据全量驻留内存——故障切换几乎瞬时,不再需要 ZooKeeper。
运行方式
4.0(2025-03)起,Kafka 彻底移除了 ZooKeeper,KRaft 成为唯一的元数据模式(不是可选项)。它的核心思路非常 Kafka:把元数据也当成一条日志。集群的所有元数据(有哪些 topic、分区分布、ISR、配置……)记录在一个内部的、单分区的 topic __cluster_metadata 里,由一组 controller 组成的 quorum 用 Raft 协议复制。
其中一个 controller 是 active controller,是元数据日志的唯一写入者;其余 controller 作为 standby 回放这条日志、把最新元数据全量驻留在内存。这样当 active controller 故障时,某个 standby 接管几乎是瞬时的——它的内存里已经是最新状态,不需要从外部系统重新加载。这正好呼应第 1 章的本质命题:连"管理集群"这件事,Kafka 都用它自己的"日志"原语来解。
| 方案 | 优势 | 为什么被取代 / 选中 |
|---|---|---|
| 外部 ZooKeeper 集群(4.0 前的模式) | 成熟的分布式协调系统、久经考验 | 是独立部署的第二套系统(多一份运维与故障面);controller 故障切换要从 ZK 全量重载元数据,分区一多重载就慢,限制了集群规模与恢复速度 |
| 每个 broker 各自持久化元数据(去中心、无 quorum) | 无单独协调组件 | 没有单一可信来源,分歧难以收敛;元数据一致性要自己从头解决——等于重造一个共识系统 |
| KRaft:元数据即日志 + controller quorum(Raft) | 无外部依赖;standby 内存常驻最新状态 → 故障切换近乎瞬时 → 支持百万级分区;元数据用 Kafka 自己的日志/复制原语 | 选中(4.0 起唯一模式) |
ZooKeeper 是一套独立系统,带来两个根本问题:多一套要运维和会故障的组件;controller 切换时要把全部元数据从 ZK 重新拉一遍、在内存重建,分区规模越大越慢。KRaft 把元数据变成一条用 Raft 复制的日志后,standby controller 持续回放、内存里始终是热的,切换不需要重载——故障恢复时间与分区数解耦,这才打开了百万级分区的天花板。
带来的代价 / 失效模式
controller quorum 多数派挂掉 = 元数据层失去可用性。Raft 靠多数派推进。3 个 controller 能容忍挂 1 个;一旦挂掉多数(3 个里挂 2 个),quorum 无法选出 active controller、无法提交元数据变更——创建 topic、leader 选举、ISR 变更等全部停摆(已有分区的纯数据读写可短时间靠缓存的元数据继续,但任何需要元数据变更的操作都会阻塞)。所以 controller 节点的部署(数量、跨机架/可用区)要按"保住多数派"来规划。
升级路径约束(截至 2026-06)。不能从 ZooKeeper 集群直接跳到 4.0——必须先在 3.x 上用迁移工具迁到 KRaft(3.9 是推荐的最后一个桥接版),再升 4.0。把"先迁 KRaft 再升大版本"当硬约束。
有人说"KRaft 把 controller 和 broker 合并了,所以集群里随便挂几台都没事"。一个 3 controller + 5 broker 的集群,如果挂掉的恰好是 2 个 controller 节点,会发生什么?
展开答案(先停 10 秒再点)
"随便挂几台都没事"是错的——要看挂的是不是 controller、有没有破坏 controller 的多数派。3 个 controller 的多数派是 2;挂掉 2 个,只剩 1 个,无法构成多数派,选不出 active controller。后果:元数据层冻结——不能建/删 topic、leader 故障后选不出新 leader、ISR 变更无法提交。已有分区如果 leader 还活着、元数据没变,数据读写可能还能撑一阵(靠各 broker 缓存的元数据),但集群已经失去自愈能力,任何一个分区 leader 再挂就无人接管。
洞察:KRaft 的可用性建立在 controller quorum 的多数派上,和 broker 数量是两回事。容量规划要分别保证"broker 副本够容忍故障"(§2.3 的 RF/ISR)和"controller 多数派够容忍故障"。把两者混为一谈是常见误区。
2.7把六个机制串起来:一条记录的一生
前六节各讲一个机制。但生产里它们是协同工作的——理解协同,才算真的理解。跟一条记录走完"producer 写入 → follower 同步 → HW 推进 → consumer 读取"这条链,看每一步踩在哪个机制上。
min.insync.replicas、HW 推进过这条记录,producer 才收到成功 ack,且consumer 才可见。持久性(向 producer 确认)和可见性(对 consumer 暴露)被 HW 这一个点同时门控。逐步拆解,注意每一步的"如果失败会怎样":
- ① producer 写(§2.2 + §2.5):记录按
murmur2(key) % 分区数选定分区,带着幂等序列号发给该分区的 leader,acks=all表示要等所有 ISR 确认。若 key 选得差 → 热分区倾斜;若没开幂等 → 重试可能乱序或重复。 - ② leader 追加(§2.1):leader 把记录顺序追加到 active segment 末尾,进 page cache,LEO +1。这一步是 O(1) 顺序写。
- ③ follower fetch(§2.3):ISR 里的 follower 各自 fetch 这条记录、写入自己的日志、推进自己的 LEO。若某 follower 掉队超过
replica.lag.time.max.ms→ 被踢出 ISR,不再参与 HW 计算。 - ④ HW 推进 + ack(§2.3):当 ISR 里所有副本的 LEO 都越过这条记录,HW 推进;此时若在线 ISR 数 ≥
min.insync.replicas,producer 收到成功 ack。若 ISR 缩到不足 → producer 收到NotEnoughReplicas、写被拒(而不是悄悄降级)。 - ⑤ consumer 读(§2.4):consumer
poll()拉取,只能看到 ≤ HW 的记录;处理完按 at-least-once 提交 offset 到__consumer_offsets。若处理超过max.poll.interval.ms没再poll()→ 被判死、分区转移、这批被重复处理。
这条链最值得记住的是 HW 是持久性和可见性的共同闸门。同一个 HW:对 producer 侧,它(配合 min.insync.replicas)决定"何时算写成功";对 consumer 侧,它决定"何时能读到"。所以"已 ack 的记录消费者要等一个 fetch 轮才看得到"不是两个独立现象,而是同一个 HW 推进过程的一体两面。把这点讲清,面试里"acks/ISR/HW/可见性"这一串就不会散架。
§本章 self-check
先合上教程,把你能想到的答案写在纸上或编辑器里。 写完再点开答案对照——直接点开等于把这一节当再读一遍。
- Kafka 为什么不在 JVM 堆里缓存消息?把缓存交给 OS page cache 换来了什么、又在哪两种情况下失效?
- 给一个带
key的 topic 加分区,为什么会"静默破坏"按 key 的顺序?为什么不能减分区? - 用一句话说清 LEO、HW、ISR 三者的关系。为什么消费者读不到 leader 上最新写入的那几条?
acks=all和min.insync.replicas各自负责什么?为什么只配acks=all不够、会在什么时候悄悄退化成acks=1?- 消费者的心跳一直正常,为什么还会被判死、导致重复处理?该调哪个参数、不该调哪个?
- 幂等 producer 的去重边界是什么?为什么"开了
enable.idempotence=true就端到端精确一次"是错的? - (跨机制综合) 一个服务"消费 topic A → 调一个外部支付 API → 生产 topic B",要求"绝不漏处理、且支付绝不重复扣款"。把 §2.3 的
acks/ISR、§2.4 的 offset 提交时机、§2.5 的事务/幂等边界串起来:Kafka 能保证哪部分?哪部分必须在 Kafka 之外解决、怎么解决?
答案(先做完再展开)
- 堆内缓存会和 OS page cache 双份存储、且几十 GB 数据进堆制造灾难性 GC。交给 page cache 换来:无 GC 压力、32GB 机器近 30GB 可用作缓存、进程重启缓存仍在内核、还能配合零拷贝直接发网卡。失效:(a) 开 TLS 时零拷贝被迫走用户态加密;(b) 滞后消费者冷读历史数据,把冷 segment 拉进缓存挤掉热数据。
- 分区由
murmur2(key) % 分区数决定,加分区改变了分母,同一个 key 的新记录可能落到与历史记录不同的分区,于是同 key 事件被劈到两个分区并行消费、顺序丢失;且无任何报错,故"静默"。不能减分区是因为减少会让已有记录的 key→分区映射失配、且要丢弃或迁移整段日志,Kafka 不支持。 - LEO 是某副本下一条要写的 offset(日志末尾);HW = ISR 中最小的 LEO(已被所有同步副本复制到的最高位);消费者只能读到 ≤ HW。读不到 leader 最新几条,是因为那几条还没被所有 ISR 复制完(在 HW 之上)——一旦 leader 挂、follower 顶上它们可能消失,所以 HW 门控不让消费者看到。
acks=all负责"等当前 ISR 所有成员确认才算写成功";min.insync.replicas负责"ISR 不许缩到这个数以下,否则拒写"。只配acks=all不够,因为 ISR 会缩——两个 follower 掉出 ISR 后 ISR 只剩 leader,"all"就是一台,等价acks=1且 producer 无感知。min.insync.replicas=2才能在这时让写入被拒。- 因为存活由调用
poll()证明,不是心跳。心跳线程在后台独立运行,但处理一批记录超过max.poll.interval.ms(默认 5 分钟)没再poll(),coordinator 就判死、转移分区,导致重复处理。该调大的是max.poll.interval.ms(或减小max.poll.records让每批更快返回);不该调session.timeout.ms——那管的是心跳/进程存活,与此无关。 - 幂等只在同一个
(PID, 分区)且同一 producer 会话内去重。producer 崩溃重启领到新 PID,broker 不认得它与旧 PID 的关系,跨重启不去重;且幂等只管 producer 到 broker 这一跳,不管"多分区写 + offset 提交"的原子性,也不覆盖外部副作用。要跨会话精确一次得用事务(transactional.id+ epoch)。 - 综合题:Kafka 这边:用
acks=all+min.insync.replicas=2+RF=3保证 topic A/B 的数据不丢;消费侧用 at-least-once(处理后再提交 offset)保证"绝不漏处理";如果只在 Kafka 内部("消费 A → 生产 B"),可以用事务把 B 的输出和 A 的 offset 提交做成原子(精确一次)。但支付 API 是外部副作用,不在 Kafka 事务范围内——Kafka 回滚撤不掉已发生的扣款。所以"支付绝不重复"必须在支付侧用幂等键(同一订单的支付请求带同一幂等 token,支付服务据此去重)解决。结论:Kafka 负责消息不丢 + 内部精确一次;"不重复扣款"靠外部幂等。把"用了 Kafka 事务就全链路精确一次"当成立是错的。
为什么 unclean leader election 能"绕过" min.insync.replicas 的保护?
你已经用 acks=all + min.insync.replicas=2 + RF=3 把写入路径锁得很死:任何时候 ISR 不足 2 就拒写,看起来已提交的记录绝不会丢。但 unclean.leader.election.enable=true 仍能让你丢已提交的数据。说清楚这条路径:为什么"写入侧的 min.insync.replicas 保护"挡不住"选举侧的 unclean 截断"?这两个机制分别在记录生命周期的哪个阶段起作用?
提示(卡住再展开)
min.insync.replicas 是写入时的闸门——它保证一条记录被 ack 时至少在 2 个副本上。unclean leader election 是故障恢复时的选择——当那 2 个(及以上)持有该记录的副本全都不可用、ISR 空了,unclean 允许一个从来没复制到这条记录的落后副本当 leader。它的日志更短,新 leader 上任后,那条"曾经满足 min.insync 地提交过"的记录在新日志里根本不存在,被当作从未发生而截断。想想 leader-epoch-checkpoint 在这里扮演什么角色,以及为什么默认值是 false——这是一道"持久性 vs 可用性,且发生在不同时间窗口"的题。