Kafka学习笔记(三):消费组、Offset与Rebalance详解
Kafka学习笔记(三):消费组、Offset与Rebalance详解
能收发消息之后,几个问题很快就会出现:
- 为什么同样的代码,换一个
group.id就又读到了旧消息? - 为什么消费者明明处理成功了,重启后还要再处理一次?
- 为什么多开几个实例,有些实例却一直没有消息?
- 为什么把
auto.offset.reset改成earliest,也没有从头读取?
这些问题的交汇点是:谁负责哪个分区,以及消费进度在哪里。
本文沿用第二篇的 order-events,共三个分区。讨论范围是 KafkaConsumer.subscribe() 的常规消费组,并以 group.protocol=classic 的示例配置解释行为。
系列导航:
一、消费组到底在管理什么?
一句话理解:
消费组把订阅的分区分配给组内消费者,并保存这组消费者对各分区的读取进度。
1.1 组内分工,组间独立
假设通知服务启动了两个实例:
1 | group.id = notification-service |
如果统计服务也需要这些事件,就使用另一个组:
1 | group.id = analytics-service |
上图只是可能的分配结果,具体分区归属取决于分配策略。稳定分配时,同一分区在一个常规消费组内由一个消费者负责;不同组则独立消费。KafkaConsumer API说明了这一机制。
1.2 为什么实例多了不一定更快?
| 分区数 | 同组消费者数 | 可能的结果 |
|---|---|---|
| 3 | 1 | 一个消费者负责三个分区 |
| 3 | 2 | 两个消费者分担三个分区 |
| 3 | 3 | 每个消费者可以负责一个分区 |
| 3 | 5 | 至少两个消费者分不到分区 |
这里假设组里只订阅这个 Topic。消费能力还取决于分区的流量分布:如果大部分事件都落在 P0,其他两个分区几乎空闲,多开实例也解决不了这个热点。
分区数限制并行度,Key 分布决定并行度能否用起来。
二、Offset 要分清三种位置
2.1 记录 Offset
这是消息在分区日志中的位置。示意:
1 | P0:[0][1][2][3][4][5][6][7] |
每个分区维护自己的位置序列,不能把 P0 与 P1 的 Offset 放在同一条数轴上比较。
2.2 当前读取位置 Position
假设一次 poll() 返回 Offset 5、6、7,此时客户端的下一读取位置可能已经到了 8。
但是,消息返回给应用,不代表应用已经处理成功。 如果应用把这三条记录扔进线程池,客户端读取位置已经推进,业务可能还没完成。
2.3 已提交位置 Committed Offset
这是保存到 Kafka 的消费进度,用于消费者恢复时确定从哪里继续。常规消费组的提交信息存储在内部 Topic __consumer_offsets 中。消费组操作文档介绍了进度管理与查询。
三个位置放在一起看:
1 | 本轮返回:Offset 5、6、7 |
这时如果提交 8,就把尚未成功处理的 7 也包含进去了。
提交位置表示“此前的记录已处理完成,恢复时可以从这个位置继续”。
以上客户端位置与已提交位置的区分,可参考KafkaConsumer 的 Offsets and Consumer Position。
三、先提交还是先处理?
3.1 先提交,后处理
1 | 取到记录 5、6、7 |
在这种消费模型下,可能出现最多一次处理:减少重复的同时,允许部分业务没有执行。
3.2 先处理,后提交
1 | 取到记录 5、6、7 |
这是常见的至少一次处理模式:用重复处理的可能性,换取崩溃后的重新尝试。前提仍包括记录没有被清理、失败能够恢复等条件。官方投递语义设计讨论了这些权衡。
实际业务通常选择:
1 | 业务持久化成功 → 提交 Offset |
四、手动提交:最容易出错的是“提交过头”
4.1 无参 commitSync() 的使用前提
第二篇采用这个流程:
1 | ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1)); |
它适用于本轮记录全部同步处理完成的情况。上面是消费循环片段,沿用第二篇的导入、客户端配置与关闭逻辑。
如果 process() 抛出异常,应该中止这一轮处理并进入恢复流程。不能捕获某条消息的错误、打印日志,再在循环结束时提交全部进度。
4.2 按分区提交完成进度
下面是替换消费循环的教学片段。process() 仍同步执行,失败时退出,不继续处理该分区的后续记录:
1 | ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1)); |
需要额外导入 java.util.Map 和 org.apache.kafka.common.TopicPartition。Kafka 4.0 的 nextOffsets() 提供本批各分区的下一位置及关联元数据,适合在该分区本批全部处理成功后提交。官方 API 示例也展示了分区级提交方式。
为什么不提交“处理成功的最大 Offset”?看这个例子:
1 | P0:100 成功,101 失败,102 成功 |
提交 103 会使 101 被跳过。只能推进到连续完成的边界,不能越过尚未成功的记录。
4.3 自动提交一定不安全吗?
也不能这样概括。问题在于提交时间能否与业务完成对齐。
同步处理完每次 poll() 的全部记录,再进行下一次轮询,自动提交也可以形成至少一次处理模式;如果记录只是入了内存队列,下一轮轮询就可能推进提交,而真正处理尚未完成。关闭消费者时也要考虑尚未完成的工作。KafkaConsumer 的手动进度说明明确指出了这些前提。
所以使用异步线程池时,不能只改成手动提交就结束了,还需要维护每个分区的连续完成进度、任务归属和停止策略。
五、auto.offset.reset:没有有效书签时怎么办?
5.1 它不是每次启动都执行
auto.offset.reset 主要用于没有初始已提交位置,或已提交位置已不在可读取范围内的情况。官方消费者配置解释了触发条件。
| 设置 | 触发重置时的行为 |
|---|---|
earliest |
从当前仍保留的最早位置开始 |
latest |
从当时日志末尾开始,等待之后的新记录 |
none |
抛出异常,交由应用处理 |
注意:earliest 不一定是 Offset 0,旧记录可能已经清理。上表列的是入门中常用的三种选项,并不是完整配置列表。
5.2 一个具体例子
1 | 当前日志保留:Offset 50~99 |
消费者从 80 继续,因为书签仍然有效。
如果这个组没有提交位置,则按 earliest 从 50 开始。如果提交位置是已经过期的 20,也会触发重置策略。
把 earliest 改一遍,不能替代主动重放。
六、Rebalance:分区的负责人变了
6.1 什么会触发重新分配?
例如消费者加入、退出或失效,订阅的分区集合发生变化,都可能让消费组重新分配工作:
1 | 最初:A ← P0、P1、P2 |
新负责人使用已提交进度恢复。如果旧负责人执行了业务但没成功提交,新的负责人可能再次执行这些记录。
6.2 两个超时不要混为一谈
| 配置 | 在 classic 协议中主要观察什么 |
|---|---|
session.timeout.ms |
组协调器多久没收到心跳后判定成员失效 |
max.poll.interval.ms |
应用调用 poll() 的最大间隔 |
max.poll.records |
一次 poll() 最多返回多少条记录给应用 |
后台心跳仍在发送,也不能无限延长业务处理时间。应用两次 poll() 间隔过长,仍可能失去分区所有权。官方消费者配置列出了这些配置的作用。
例如,假设每条事件同步处理耗时约 1 秒,一批 500 条就可能花 500 秒。不能只看到“消费者还有心跳”,就忽略轮询间隔。
排查时优先测量处理耗时与慢调用,再调整批量大小。max.poll.records 限制本轮返回量,并不等于限制底层每次网络获取或客户端缓存的全部数据量。
6.3 Kafka 4.0 的新消费组协议
Kafka 4.0 提供新的 Consumer Rebalance Protocol,可通过 group.protocol=consumer 选择。它使用服务端管理分配,并改善重新分配的方式。官方协议说明介绍了迁移步骤。
新协议下心跳与会话超时由对应 Broker 配置管理,不能直接套用 classic 的所有调参经验。第二篇明确设置 classic,是为了让这组实验保持统一的解释范围。
6.4 回调解决不了所有并发问题
ConsumerRebalanceListener 可以在分区撤销、分配、丢失时配合应用清理状态。需要维护的是:
- 已完成记录对应的提交边界。
- 已撤销分区还有哪些任务正在执行。
- 分区所有权丢失后,旧任务如何停止或限制写入。
回调中的提交也可能失败。生产环境采用线程池消费时,需要完整的分区任务管理与业务幂等,不能仅在回调中补一个 commitSync()。
七、观察积压与分区分配
7.1 看消费进度
1 | docker exec kafka-learning /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group order-java-demo-v1 |
在普通日志下,可以先把某分区的 LAG 理解成:
1 | LAG = LOG-END-OFFSET - CURRENT-OFFSET |
但它只是 Offset 差。事务记录、日志压缩造成的间隙等,都会让“差值”与业务消息条数不完全一致。
7.2 看分区负责人
1 | docker exec kafka-learning /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group order-java-demo-v1 --members --verbose |
查询与进度重置的操作可以参考官方消费组管理文档。
7.3 积压先找原因
| 观察结果 | 优先排查 |
|---|---|
| 所有分区持续落后 | 总消费能力、下游慢调用、资源限制 |
| 只有一个分区严重落后 | Key 热点、该分区的异常事件 |
| 进度反复停住 | 处理失败、提交失败、反复 Rebalance |
| LAG 较低但业务仍有延迟 | 客户端缓存、异步任务队列、下游排队 |
估算追赶时间时可以使用:
1 | 追赶时间 ≈ 当前积压 /(消费速率 - 新增速率) |
这个估算要求消费速率高于新增速率,且消息处理成本相对稳定。如果消费与新增都是每秒 1000 条,积压不会自行减少。
八、做三个小实验
实验一:观察组内分配
同时运行两个第二篇的 Java 消费者,它们使用同一个 group.id。用 --members --verbose 查看分区归属,再关闭其中一个实例,观察剩下的实例接手分区。
注意:Java 生产者只发送 O1001,消息会集中在一个分区,所以“有分区但没打印消息”并不能证明没有分配成功。可用命令行发送多个订单 Key,再对照查看。
实验二:比较相同组与新组
处理消息并成功提交后,重启相同组,观察它继续等待新消息。然后将 Java 代码中的组名改成一个从未用过的名称,保留 earliest,观察它读取保留中的旧消息。
实验三:预览重放范围
先停止该组的所有消费者,再执行预览:
1 | docker exec kafka-learning /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group order-java-demo-v1 --topic order-events --reset-offsets --to-earliest --dry-run |
确认这个练习组的历史事件可以再次处理后,再实际重置:
1 | docker exec kafka-learning /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group order-java-demo-v1 --topic order-events --reset-offsets --to-earliest --execute |
这里只针对本地练习组。真实业务重放时,要先明确是否会重新发短信、重复扣库存,以及处理结果如何去重。
Offset 是消费恢复的依据,业务幂等是重复处理时的保护。 下一篇把这两层与生产者、副本、Kafka 事务串起来。