Kafka 底层:分区日志、复制与消费位点

沿订单事件从发送到业务提交的路径,解释 Kafka 分区日志、批次与稀疏索引、ISR 与复制确认、HW 与 LSO、消费组位点和 KRaft 的边界。

一条订单事件可以被库存服务、数据分析和风控分别读取。库存服务重启以后,还可以从上次提交的位置继续读取;分析服务落后几个小时,也不必要求生产者重新发送。Kafka 怎样用同一份日志支持这些不同的消费进度?当机器故障、消费者换人或者消息被清理以后,这套能力又有哪些限制?

本文以 Kafka 4.1 文档和 4.1.0 存储源码为依据,讨论本地分区日志、常规消费组与 KRaft。分层存储和 Share Group 不在本文范围内;“一个分区在组内分配给一个消费者”专指常规消费组。涉及索引时直接对照源码,避免把逻辑 offset 误写成文件字节地址。

前面的文章已经讨论了 消息不丢失 和 重复消费与幂等。这里把它们下面的机制串起来:记录先进入某个分区的日志,复制决定它能否安全可见,消费者自己的进度决定故障后从哪里继续,外部业务事务则决定订单究竟发生了什么。

一、分区日志为什么适合多组消费者

Topic 是业务层的名称,Partition 是实际组织记录、分配读取任务和进行复制的单位。一个 Topic 有三个分区,就有三条相互独立的有序日志。订单事件落在分区 0 的 offset 100,支付事件落在分区 1 的 offset 100,两者没有共同的“第 100 条”含义,也不能靠数字判断先后。

每个分区副本都保存相应日志。正常写入由该分区的 Leader 接收,其他副本作为 Follower 跟随;不同分区的 Leader 可以分散到不同 Broker 上。因此一个 Broker 可能同时承担一些分区的 Leader 和另一些分区的 Follower,不是整个集群只选出一台机器负责全部订单。

消费者读完消息,并不会直接从日志中把它删除。库存消费组与风控消费组分别保存进度,可以在同一条日志上前后移动。这个模型适合事件被多个下游独立使用:新增一个分析消费者,通常不需要生产者为它额外保存一份专用消息。

业务需要同一订单有序时,通常以订单 ID 作为分区路由的 key,并保持路由规则稳定。不过分区只能保存收到的次序,不能修正生产者本来就发反的业务事件。扩分区可能改变 key 的落点;消费者再把同一分区的消息无约束地扔进线程池,也会打乱实际执行次序。日志有序和业务有序之间,仍然隔着发送与执行协议。

分区数也是并行度与管理成本之间的选择。常规组内,一个分区在正常分配状态下由一个消费者负责,增加消费者不能无限拆分同一个热点分区。若一个大客户占据绝大多数流量,只扩机器不改路由,瓶颈可能仍停留在那一条日志和它的业务处理器上。

二、一条日志为什么还要切成多个段

把一个分区长期写成单个巨大文件,会让过期清理、恢复和定位都难以管理。Kafka 将日志分成 Segment,也就是日志段。当前活跃段承接追加,达到相应大小或时间条件后滚动,后续写入进入新段。段的起始 offset 用来标识它在分区日志中的范围。

正文保存在 .log 文件中,配套索引帮助定位。这里最重要的是 offset 索引:它把分区内的逻辑位置映射到相应段中的物理位置。时间索引辅助按时间寻找候选位置,事务索引参与识别相关事务范围。它们服务于不同查询,不能把所有索引都叫成“消息 ID 到正文”的映射。

逻辑 offset 表示记录在分区日志中的位置;物理 position 表示段文件内的字节位置。前者跨段延续,后者进入新段后重新从文件起点计量。某条记录正文变大,会影响后续字节地址,却不意味着 offset 必须按正文长度增长。

例如一个段以 offset 100 开始,另一个以 1000 开始。查 105 时先选择前一个段;查 1200 时选择后一个段。段内的相对 offset 可以表示为“目标位置减去段起点”,索引不必为每个条目都重复保存完整的长整型起点。

这种组织使删除旧数据时能够按段处理,而不必不断移动仍保留的消息。代价是保留策略和删除时点存在段粒度与后台调度边界:配置保留一天,不表示第 86400 秒每条记录都准确消失。容量规划要考虑活跃段、滚动条件及清理延迟,不能只拿平均每日流量乘一个天数。

具体段内结构可以从 LogSegment 中查看。它同时持有正文和相应索引,读取时先定位范围,再进入正文。

三、稀疏索引怎样找到目标 offset

Kafka 的 offset 索引不是每条记录一个条目。它间隔记录一些映射,所以叫稀疏索引。查询时先用索引缩小范围,再从候选字节位置顺着日志找目标。它用一小段扫描换取更小的索引和更少的索引更新。

假设段起点是 100,索引中存在 102 → 300 和 110 → 1200 两个条目,目标是 105。105 没有索引项,不代表消息不存在。查询会找到不大于 105 的最近条目 102,从文件的 300 字节附近进入,然后向后查找包含目标位置的记录批次。

稀疏索引定位后扫描正文

这里的“扫描”有明确范围,不能想象成每次都从分区第一条记录读到最后。先选段,再查索引,再扫描局部正文。索引更密,局部扫描可能更少,但索引体积和维护开销会上升;索引更疏,映射更少,查询需要读过更多中间数据。这个选型与按任意字段建立数据库索引的目标不同。

4.1.0 的 OffsetIndex 使用相对 offset 与物理位置组织条目;lookup 寻找不超过目标的映射。随后 LogSegment.translateOffset 调用正文查找逻辑,继续定位记录批次。本文示意中的数值用于解释关系,不代表默认索引间隔。

另一个容易误解的地方是 offset 不保证消费者看到的记录逐个连续。日志压缩可能移除中间记录,事务控制记录也占据日志位置,终止的事务记录不会交给 read_committed 应用。读到 105 后再读到 108,不能单凭缺了 106、107 就判定丢失。排查需要结合清理策略、隔离级别和原始日志范围。

四、批次怎样连接网络与存储

生产者通常不会为每条业务记录都单独完成一次网络往返。客户端按分区聚合记录,形成批次,再发送给相应 Leader。批次能摊薄请求头、网络往返和追加操作的固定开销,也让压缩更容易利用相邻记录中的重复内容。

batch.size 与 linger.ms 分别影响批次容量和等待聚合的时间,但不是“每条消息必须等待这么久”。批次已经满足发送条件时可以提前发送;Broker 背压、网络与重试又可能让实际延迟超过这段聚合等待。低流量场景和满负载场景,不能只用同一个吞吐数字描述。

可以把一个 Record Batch 简化理解成下面的形状。它展示的是概念字段,不是可直接序列化的协议代码:

Record Batch
  起始 offset、最后一条记录的 offset 差值
  长度、格式版本、校验信息、压缩与事务标记
  producerId、producerEpoch、起始序列号
  一组 records:key、value、timestamp、headers 等

批次头不仅帮助传输,也承载校验、幂等和事务所需的元信息。正文日志保存批次,读取侧也按批次处理,因此索引定位的是进入日志的候选位置,不能理解成按 offset 直接跳到一个固定长度 JSON 对象。

压缩会减少网络和存储字节,但消耗 CPU。字段高度重复的埋点批次,与已经压缩过的图片片段,收益不会相同。调参需要同时看生产等待、请求时延、压缩比、Broker CPU 和消费者解压成本。只把批次做大,可能让吞吐改善而尾延迟变差。

批次格式依据 Kafka Message Format,聚合配置依据 Producer Configs。生产者侧幂等解决的是协议内重试重复写入;应用再次调用发送同一业务事件,不会因为 value 一样就自动被判成重复。

五、页缓存与零拷贝各解决什么问题

Kafka 使用文件系统保存日志,并利用操作系统页缓存。写入文件与数据真正刷入持久介质之间有时间差;读取近期热点日志,也可能直接命中页缓存,不需要每次重新访问磁盘。因此“消息在文件中”与“每条发送都同步落盘”不是一个结论。

从工程角度看,批量追加、尽量连续的访问和缓存复用共同减少开销。它们不能让硬件上限消失:多个分区同时写入时,会有多个文件尾部活跃;历史回放落到冷数据时,需要实际读取存储;大量旧数据读取还可能挤占近期消息的缓存。观察一个追赶中的消费者,和观察一个回放半年前日志的消费者,看到的磁盘压力可能完全不同。

所谓零拷贝,主要指合适路径下减少正文经过用户态缓冲区的重复复制。文件到网络传输可以利用操作系统支持的优化,并非数据没有经过内存、网卡也不用搬运字节。Kafka 4.1 的 设计文档 还明确指出,启用 SSL 时不能直接沿用其中描述的 sendfile 路径,因为加密处理需要相应用户态路径。

所以不能用一句“Kafka 快是因为零拷贝”解释所有部署。请求批次、数据冷热、TLS、压缩、复制流量和消费者处理速度都在实际路径上。优化时应先确认瓶颈是网络、存储、CPU,还是应用没有及时拉取;再决定调整哪一层。

同样,不宜为了追求安全就机械地把每条记录强制刷盘。同步刷盘有延迟与吞吐代价,复制也有自己的故障覆盖范围。应用需要先定义允许承受哪些故障,再选择刷盘与复制策略,相关可靠性问题在 消息怎样不丢失 中有完整讨论。

六、复制确认为什么不是简单的多数票

每个分区都有自己的副本集合,Leader 接收追加,Follower 拉取并追赶。ISR 是 In-Sync Replicas,即被认为同步状态满足要求的副本集合,包含 Leader。副本落后到不满足要求时可能被移出,追上之后才重新加入;它不是永远等于配置的复制因子。

假设复制因子为 3,分别是 A、B、C,A 为 Leader。此时 ISR 可以是三台,也可以因 C 落后变成 A、B。acks=1 关注 Leader 的本地追加确认,acks=all 则等待当时 ISR 的确认条件。发送超时仍可能发生在数据已经写入、响应没有回来的时刻,客户端不能据此认定未写入。

min.insync.replicas 为这套确认设置最低同步副本门槛。若配置 2,且使用 acks=all,ISR 为三台时并不是任选两台返回就算完成;ISR 缩到两台时需要这两台满足确认;只剩一台时则不满足最低门槛。这是牺牲部分写入可用性来保留所需的安全边界。

# 解释三副本、至少两份同步副本的组合;不是通用调参模板
# Topic:replication.factor=3, min.insync.replicas=2
# Producer:
acks=all
enable.idempotence=true

上面两条生产者配置也不等于“所有故障下绝不丢”。副本是否处于独立故障域、是否允许不安全选主、集群是否同时断电,以及相关客户端配置是否兼容,都影响结果。复制确认也不证明所有副本已经对每批消息执行了物理刷盘。

4.1 还需要考虑 ELR,即 Eligible Leader Replicas。它记录特定条件下仍有资格接任 Leader 的副本,扩展了单纯“只看当前 ISR”的选主理解。新建 4.1 集群默认启用这一能力,并强化最低 ISR 对高水位推进的约束;不能把任意落后的副本都当成安全候选。版本边界参见 ELR 文档,确认参数参见 Topic Configs。

七、LEO、HW 与 LSO 为什么要分开

消息已经出现在 Leader 的本地文件里,不意味着普通消费者马上能读到它。至少要区分本地日志末尾、复制安全的可见范围,以及事务隔离允许交付的范围。它们分别对应 LEO、HW 和事务语境下的 LSO。

LEO 是 Log End Offset,表示下一条追加的位置。HW 是 High Watermark,高水位,表示复制提交的边界。本文所有边界都使用右开区间:HW 为 110 时,普通读取最多涉及小于 110 的位置,不是已经包括 offset 110。

LSO 在这里指 Last Stable Offset。对 read_committed 来说,还要避免越过尚未结束的事务,所以稳定读取边界不能高于 HW,也可能被较早的未完成事务挡住。这里的 LSO 不要与日志起始位置 Log Start Offset 混淆,后者描述最早仍可读取的数据位置。

本地追加、复制提交与事务稳定边界

用图中的数值解释:LEO 是 120,HW 是 110,LSO 是 100。110 到 119 的记录虽然在 Leader 本地,却没有进入相应复制可见范围。100 到 109 已经过了复制边界,但事务稳定读取暂时不能越过 100。read_uncommitted 也不能直接读取所有本地未复制尾部,“未提交”在这个配置名中主要指事务隔离。

LSO 以下也不意味着每条日志记录都交付给业务。控制记录和已终止事务中的记录仍需被过滤。一个长时间未结束的事务,可能挡住后面已经完成的事务记录,这是排查“日志有数据但消费者读不到”时需要检查的另一条路径。

这三个位置都不代表库存服务处理进度。消费者可能只处理到 80,也可能已经将业务处理到 99、尚未提交消费位点。把 Broker 的可见性边界与应用进度分开,才能判断等待发生在复制、事务还是业务侧。相关读取语义见 KafkaConsumer API。

八、消费组位点究竟保存什么

常规消费组有一个组标识,同一组的成员分担分区,不同组独立推进。两个系统若不小心配置相同的 group.id,可能变成互相分担消息,而非各自完整收到消息。这个问题不是 Topic 数据丢了,而是业务期望与分组配置不一致。

消费者本地 position 是下一次交付的位置,随着 poll 返回记录向前推进。committed offset 则是持久保存的恢复位置。业务处理完成的位置又是第三种进度:poll 已返回一批记录,不能证明这些记录对应的数据库事务已经提交。

例如拉取 100 到 109 后,本地 position 可以已经是 110,而应用只完成到 103。如果此时把 110 提交成恢复位点,重启后可能直接跳过尚未完成的 104 到 109。如果全部业务完成后尚未来得及提交,重启会重复读取这一批,则需要业务幂等来接住。

提交值通常表示下次应读取的位置。连续完成到 109,提交 110;不是提交 109 来表示“最后成功的那一条”。存在 offset 空洞时,也不能机械地按成功条数加一计算,需要按实际返回位置与完整处理范围管理。

位点提交交给消费组协调器,记录持久化在内部 Topic __consumer_offsets,协调器还会维护缓存以便恢复查询。它与业务 Topic 一样需要可靠的日志与复制,不过内容是组进度等内部状态。结构可以参见 Consumer Offset Tracking。

因此 Kafka 并没有替每条业务消息保存一份“库存已扣减、短信已发送”的完整清单。它记录的是组与分区的恢复位置。应用把一个位置提交上去,相当于告诉恢复流程可以从那里接着做;这个判断的真实性由应用负责。

九、并发消费为什么容易把进度提交错

应用想提升吞吐时,常见做法是一个线程 poll,再把记录交给工作线程。问题在于不同任务完成速度不同:offset 100 仍在等待数据库,101、102 已经完成。如果提交 103,故障后 100 就可能被跳过。

安全推进需要维护已经连续完成的前缀。后面的任务完成可以先记录下来,但不能跨过前面的未完成任务提交进度。不同分区可以分别计算,避免一个分区的慢任务拖住全部分区;同一分区需要业务有序时,则还应限制执行并发,不能只解决提交位点。

业务完成与位点提交之间的故障窗口

图中先提交数据库,再提交消费位点。两步中间崩溃时,新消费者会再次读到消息,因此应以事件 ID 或业务操作键,在同一个业务事务中完成去重与状态变更。若先提交位点再做业务,故障窗口就变成消息被跳过,通常更难恢复。

生产者幂等无法覆盖这个窗口。它避免协议内重试造成重复追加,无法阻止消费端对同一个外部接口调用两次。Kafka 事务可以将 Kafka 输出与相应输入位点放进一项事务,配合下游 read_committed 构建相应处理语义;普通 MySQL 更新、发券与扣款并不会自动加入其中。

线程模型也有边界:Java KafkaConsumer 不是可被任意多个线程并发调用的对象。工作线程可以处理业务结果,再把进度交回负责消费者的线程;不能为了省事让每个工作线程同时调用提交或轮询方法。并发模型、背压和组成员失效处理需要一起设计,而不是仅把自动提交关掉就结束。

具体业务去重方案在 重复消费与幂等 中展开;这里重点是底层位点模型为什么要求它存在。

十、Rebalance 换的是分区责任,不会撤销外部业务

消费者加入、离开或被判断失效时,消费组需要重新分配分区。新的成员根据相应恢复进度接手。原成员如果已拿到一批记录但还没提交位点,新成员就可能再次处理它们。

更危险的是旧成员的业务任务没有立刻停止。它的网络可能只是暂时卡住,工作线程仍在执行数据库操作;新的成员已经接手相同分区。组协议可以拒绝某些过期的协议操作,但不能进入你的数据库把已经发出的业务请求撤销。外部副作用需要幂等、状态条件或业务侧有效的 fencing 约束。

分区撤销时,应用应停止为相应分区继续派发任务,处理允许完成的在途结果,并且只提交真正完成的连续范围。若任务无法安全完成,宁可让新成员重读,也不要把未知结果宣称为已完成。停止过程必须有时间边界,否则恢复会被长期悬挂任务拖住。

Kafka 4.0 起的新消费组协议已经可用,4.1 客户端需要通过 group.protocol=consumer 选择它,不能因为 Broker 支持就认为所有应用自动切换。该协议引入服务器端分配与更增量的协调,减少全组同步障碍;Classic 协议仍然存在。心跳和超时参数也有不同管理方式,升级时应按照 Consumer Rebalance Protocol 核对。

协议改进可以减少分配波动与暂停,无法替代业务安全设计。观察重平衡时,需要同时看组成员变更、分区归属、撤销耗时和重复业务命中率;只看到消费者数量恢复正常,不等于旧任务和新任务没有重叠。

十一、保留与压缩为什么决定可回放的边界

消费位点可以回退,但只能回到仍存在的数据范围。cleanup.policy=delete 按相关保留条件清理旧日志段,并不等待每个消费组都读完。下游停机太久,位点可能落到最早保留位置之前;日志已经删除,重新设一个旧 offset 也变不回来。

这时自动重置策略只决定从哪个有效位置开始,不能修复缺失的业务结果。对订单和账务数据,应先判断缺口是否需要通过源数据库、归档或对账恢复,再决定如何重置。未经检查就设成 latest,可能让应用看似恢复运行,却永久跳过停机期间的一段业务。

cleanup.policy=compact 则用于保留各个 key 的较新状态,后台逐步清理较旧记录。它适合以最新值恢复状态的场景,不适合把每一次修改都当成永久审计历史。清理是异步的,日志中暂时存在同一 key 的多个版本完全正常,也不能把压缩当成实时重复消息去重。

null value 的删除标记还有保留窗口;落后的恢复读取必须考虑能否及时看到这些删除。一个 key 曾经被删除,消费者若错过标记,仅靠过旧的本地状态可能继续保留它。组合 delete 与 compact 时,也应理解两种清理条件同时影响历史范围。

这些配置依据 Topic Configs。从业务角度,保留窗口应覆盖预期停机、补数与回放时长,并留出恢复追赶的余量。容量不足与无限保留之间没有免费的选择:数据可靠性还需要明确归档职责和恢复入口。

十二、KRaft 管的是元数据,不是每条订单多数投票

Broker 需要知道哪些分区存在、有哪些副本、当前 Leader 是谁,以及有关配置。Kafka 4.1 的 KRaft 使用控制器仲裁维护集群元数据;控制器有活动与备用角色,元数据仲裁需要相应多数可用。Broker 与 Controller 可以配置为不同进程角色。

这套控制面的 Raft,与业务分区的数据复制不是同一个协议实例。订单写入仍进入分区 Leader,并沿该分区副本和 ISR 等机制推进;不能把 acks=all 解释成“向 KRaft 控制器发消息,三台里两台投票通过”。把两条路径混在一起,会同时误解写入确认和选主恢复。

控制器失去足够仲裁能力,元数据变更和故障恢复会受影响,但不能简单描述成某个控制器一停,所有现存分区立刻从磁盘消失。反过来,控制器仲裁正常,也不代表所有业务分区副本健康。控制面与数据面应分别监控。

消费组协调器同样不是控制器的另一个名字。它是 Broker 上承担相应组协议与进度管理的角色,组数据进入内部 Topic;应用发券事务则在自己的数据库。分清这几种状态的持有者,才知道每次报错应该去查哪一侧。

部署角色和仲裁边界参见 KRaft 文档。本文不展开元数据日志协议细节,因为它对理解订单读取的第一条主线不是必需;但必须保留这一边界,避免把新版 Kafka 又套回旧版 ZooKeeper 的角色图。

十三、沿一条订单事件检查整条链路

假设订单服务发出事件 order-paid-42,key 为订单 ID,进入订单 Topic 的分区 2。客户端将它与相同分区的一组记录组成批次,发给当前 Leader。Leader 将记录追加到活跃段,并维护相应索引,事件获得该分区中的 offset 105。

Follower 追赶这段新增日志。在三副本、最低 ISR 为 2、acks=all 的设定下,需要符合相应复制确认条件才能给生产者成功响应。若响应丢失,生产者重试需要依靠协议幂等;若订单服务重新构造一次业务发送,则仍应带稳定事件标识,让下游能识别业务重复。

当记录进入复制可见范围,并满足消费者使用的事务隔离条件后,库存消费组从自己的进度读取。存储先按 offset 找段与索引,再读取批次。分析消费组可能仍在处理 80,并不妨碍库存组读取 105;两者也不会因为其中一组完成就把正文立即删除。

库存应用在数据库事务中记录事件去重键并推进库存业务状态。提交完成以后,消费进度才允许越过 105。如果此时进程退出,数据库已经成功而位点尚未成功,新成员重读事件,去重约束阻止重复生效,再推进消费位点。两个系统之间没有凭空出现原子事务,安全性来自可重试的业务协议。

排障也可以顺着这条链路走。生产者迟迟没有成功响应,检查发送重试、ISR 与请求延迟;日志末尾增长但 HW 不前进,检查副本与相关门槛;HW 前进而事务读取等待,检查未完成事务;应用已拉取却不推进,检查任务背压与数据库;位点快速前进而业务缺结果,重点检查是否提前提交。

此外,要区分“读取落后”和“业务落后”。监控若只看客户端已读取的位置,工作线程队列积压可能被藏起来;监控若只看提交位点,也应结合处理耗时和外部状态判断瓶颈。关键指标应带 Topic、Partition、Group 等作用域,并与事件 ID、业务对象和重试记录关联,才方便解释某一笔订单的实际结果。

十四、与 RocketMQ 对比时应该比较什么

上一篇 RocketMQ 存储架构 固定分析 4.9.8 经典实现。把它和本文 Kafka 4.1 的本地存储对照,最有价值的是比较正文组织与定位方式,而不是仅比较名词数量。

比较点 RocketMQ 4.9.8 经典存储 Kafka 4.1 本地分区日志
正文组织 Broker 内 CommitLog 混合追加各队列正文 每个分区副本有自己的分段日志
按消费位置定位 ConsumeQueue 条目定位 CommitLog 正文 选日志段,查稀疏索引,再扫描批次
索引关系 经典队列中每条消息对应定长队列条目 offset 索引不为每条记录保存映射
消费后清理 消费完成不等于立即删除正文 位点提交不等于立即删除日志

共同点是利用追加日志承载消息,并把存储位置与业务完成分开。差异则会影响文件管理、索引写入和读取路径,但不能单凭“集中追加”或“分区追加”就判断所有负载谁更快。分区数量、消息大小、数据冷热、复制模式和过滤需求都会改变实际结果。

对业务选型,更值得先问:是否需要多个独立系统回放同一事件流,回放窗口多长,是否围绕流式处理生态工作,任务失败怎样恢复。顺序、延迟、幂等这些需求仍然要分别落实,底层日志不会自动替应用选择一套业务协议。

理解 Kafka 时,我会把三个问题始终分开:数据现在保存在什么范围,协议允许消费者看到什么范围,业务已经可靠完成了什么范围。分区日志、复制边界和消费位点各回答其中一部分;只有把外部业务提交也放进链路,才能说明一次订单事件究竟有没有安全生效。