Kafka学习笔记(三):消费组、Offset与Rebalance详解

能收发消息之后,几个问题很快就会出现:

  • 为什么同样的代码,换一个 group.id 就又读到了旧消息?
  • 为什么消费者明明处理成功了,重启后还要再处理一次?
  • 为什么多开几个实例,有些实例却一直没有消息?
  • 为什么把 auto.offset.reset 改成 earliest,也没有从头读取?

这些问题的交汇点是:谁负责哪个分区,以及消费进度在哪里。

本文沿用第二篇的 order-events,共三个分区。讨论范围是 KafkaConsumer.subscribe() 的常规消费组,并以 group.protocol=classic 的示例配置解释行为。

系列导航:

  1. 核心概念与架构
  2. Docker 与 Java 入门实战
  3. 消费组、Offset 与 Rebalance
  4. 消息可靠性、幂等与事务
  5. 面试高频问题总结

一、消费组到底在管理什么?

一句话理解:

消费组把订阅的分区分配给组内消费者,并保存这组消费者对各分区的读取进度。

1.1 组内分工,组间独立

假设通知服务启动了两个实例:

1
2
3
4
group.id = notification-service

Consumer A ← P0、P1
Consumer B ← P2

如果统计服务也需要这些事件,就使用另一个组:

1
2
3
group.id = analytics-service

Consumer C ← P0、P1、P2

上图只是可能的分配结果,具体分区归属取决于分配策略。稳定分配时,同一分区在一个常规消费组内由一个消费者负责;不同组则独立消费。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
2
3
4
本轮返回:Offset 5、6、7
当前 Position:8
业务完成情况:5、6 成功,7 仍在处理中
Committed Offset:5(本轮开始前的进度)

这时如果提交 8,就把尚未成功处理的 7 也包含进去了。

提交位置表示“此前的记录已处理完成,恢复时可以从这个位置继续”。

以上客户端位置与已提交位置的区分,可参考KafkaConsumer 的 Offsets and Consumer Position。


三、先提交还是先处理?

3.1 先提交,后处理

1
2
3
4
5
6
7
8
9
取到记录 5、6、7
↓
提交下一位置 8
↓
开始业务处理
↓
处理到 6 时崩溃
↓
重启从 8 继续:6、7 没有完成,却被进度跳过

在这种消费模型下,可能出现最多一次处理:减少重复的同时,允许部分业务没有执行。

3.2 先处理,后提交

1
2
3
4
5
6
7
8
9
取到记录 5、6、7
↓
全部处理成功
↓
提交 8 之前崩溃
↓
重启仍从此前提交的 5 继续
↓
5、6、7 可能再次执行

这是常见的至少一次处理模式:用重复处理的可能性,换取崩溃后的重新尝试。前提仍包括记录没有被清理、失败能够恢复等条件。官方投递语义设计讨论了这些权衡。

实际业务通常选择:

1
2
3
业务持久化成功 → 提交 Offset
+
重复到达时用业务幂等保护结果

四、手动提交:最容易出错的是“提交过头”

4.1 无参 commitSync() 的使用前提

第二篇采用这个流程:

1
2
3
4
5
6
7
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
for (ConsumerRecord<String, String> record : records) {
process(record); // 同步执行,失败直接向外抛异常
}
if (!records.isEmpty()) {
consumer.commitSync();
}

它适用于本轮记录全部同步处理完成的情况。上面是消费循环片段,沿用第二篇的导入、客户端配置与关闭逻辑。

如果 process() 抛出异常,应该中止这一轮处理并进入恢复流程。不能捕获某条消息的错误、打印日志,再在循环结束时提交全部进度。

4.2 按分区提交完成进度

下面是替换消费循环的教学片段。process() 仍同步执行,失败时退出,不继续处理该分区的后续记录:

1
2
3
4
5
6
7
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
for (TopicPartition partition : records.partitions()) {
for (ConsumerRecord<String, String> record : records.records(partition)) {
process(record);
}
consumer.commitSync(Map.of(partition, records.nextOffsets().get(partition)));
}

需要额外导入 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
2
3
当前日志保留:Offset 50~99
已有提交位置:80
配置:earliest

消费者从 80 继续,因为书签仍然有效。

如果这个组没有提交位置,则按 earliest 从 50 开始。如果提交位置是已经过期的 20,也会触发重置策略。

把 earliest 改一遍,不能替代主动重放。


六、Rebalance:分区的负责人变了

6.1 什么会触发重新分配?

例如消费者加入、退出或失效,订阅的分区集合发生变化,都可能让消费组重新分配工作:

1
2
3
4
5
6
7
8
最初:A ← P0、P1、P2

B 加入
↓
重新分配
↓
A ← P0、P1
B ← 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 事务串起来。