Consumer Group(消费组)
定义与作用
Consumer Group 是 Kafka 消费模型的核心抽象。多个 Consumer 组成一个消费组,共享同一个 group.id,协同消费一个或多个 Topic。消费组解决了 Kafka 消费端的水平扩展问题:增加 Consumer 实例即可提升消费吞吐。
消费组的关键规则:一个 Partition 在同一个消费组内最多只能被一个 Consumer 消费。
核心原理
消费组模型
核心推论:
| 场景 | 结果 |
|---|---|
| Consumer 数量 = 分区数 | 每个 Consumer 负责 1 个分区(理想状态) |
| Consumer 数量 > 分区数 | 多余 Consumer 空闲,浪费资源 |
| Consumer 数量 < 分区数 | 部分 Consumer 负责多个分区 |
Rebalance(重平衡)
Rebalance 是消费组内分区所有权重新分配的过程。
Rebalance 触发条件:
| 触发条件 | 说明 |
|---|---|
| Consumer 加入 | 新 Consumer 加入消费组 |
| Consumer 离开 | 正常关闭或心跳超时 |
| Consumer 崩溃 | session.timeout.ms 内未发送心跳 |
| 分区数变更 | Topic 分区数增加 |
| 订阅的 Topic 变更 | 使用正则订阅匹配到新 Topic |
Rebalance 的代价:在 Rebalance 期间,整个消费组暂停消费(Stop-The-World),直到分区重新分配完成。这是 Kafka 消费端最大的性能抖动来源。
完整示例
示例一:多 Consumer 实例水平扩展
场景:Topic logs 有 6 个分区,从 1 个 Consumer 逐步扩展到 3 个。
操作前环境:Topic logs,6 分区。Consumer Group log-processor 初始 1 个 Consumer。
步骤:
# 终端 1:启动 Consumer 1
#(Java 代码)consumer.subscribe("logs")
# 查看消费组状态
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group log-processor --describe
# GROUP TOPIC PARTITION CURRENT-OFFSET LAG CLIENT-ID
# log-processor logs 0 10000 0 consumer-1
# log-processor logs 1 12000 0 consumer-1
# log-processor logs 2 11000 0 consumer-1
# log-processor logs 3 10500 0 consumer-1
# log-processor logs 4 11500 0 consumer-1
# log-processor logs 5 11800 0 consumer-1
# → 1 个 Consumer 处理 6 个分区
# 终端 2:启动 Consumer 2(相同 group.id)
# → Consumer 1 的 3 个分区被撤销,分配给 Consumer 2
# 终端 3:启动 Consumer 3
# → 每个 Consumer 各负责 2 个分区
操作后状态:
# 3 个 Consumer 协同消费
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group log-processor --describe
# log-processor logs 0 10000 0 consumer-1
# log-processor logs 1 12000 0 consumer-1
# log-processor logs 2 11000 0 consumer-2
# log-processor logs 3 10500 0 consumer-2
# log-processor logs 4 11500 0 consumer-3
# log-processor logs 5 11800 0 consumer-3
操作前后对比:
| Consumer 数 | 每 Consumer 分区数 | 总消费吞吐 | 说明 |
|---|---|---|---|
| 1 | 6 | ~60 MB/s (假设单分区 10MB/s) | 单机瓶颈 |
| 2 | 3 | ~100 MB/s | 提升近 2x |
| 3 | 2 | ~120 MB/s | 理想状态(分区数=Consumer数) |
| 4 | 1 空闲,3 各 2 | ~120 MB/s | 多余 Consumer 浪费 |
示例二:独立消费组实现多路消费
场景:同一个 orders Topic,一个消费组做实时处理,另一个消费组做数据归档。
# 实时处理消费组
# Group: realtime-processor, Offset 从 latest 开始
bin/kafka-console-consumer.sh --topic orders \
--group realtime-processor \
--bootstrap-server localhost:9092
# 数据归档消费组(独立,不受上面影响)
# Group: data-archiver, Offset 从 earliest 开始
bin/kafka-console-consumer.sh --topic orders \
--group data-archiver \
--from-beginning \
--bootstrap-server localhost:9092
操作后验证:
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
# realtime-processor
# data-archiver
# 两个消费组完全独立,各自维护 Offset
易错场景
易错 1:Consumer 数 > 分区数
场景:Topic 有 3 个分区,部署了 5 个 Consumer 实例。
后果:2 个 Consumer 处于空闲状态(不消费任何分区)。不仅浪费资源,每次 Rebalance 开销也更大。
规则:max_parallelism = min(partition_count, consumer_count)。
易错 2:Rebalance 期间的消息堆积
场景:消费组频繁触发 Rebalance(如 Consumer 健康检查过于激进)。
后果:每次 Rebalance 暂停消费 5-30 秒,高频 Rebalance 导致消息严重堆积,而且消费者永远追不上。
解决:
- 增大
session.timeout.ms和max.poll.interval.ms - 使用 Cooperative Rebalance(增量重平衡,减少暂停时间)
- 设置
group.instance.id使 Consumer 成为静态成员
面试高频考点
Q:同一个消费组内,如果 Consumer A 崩溃,它的分区会怎样?
A:Broker 检测到 Consumer A 未在 session.timeout.ms(默认 45s)内发送心跳后:
- 触发 Rebalance
- 撤销消费组内所有 Consumer 的分区分配
- 重新分配——Consumer B 可能继承 Consumer A 原来的分区
- Consumer B 从该分区的已提交 Offset 继续消费
此时可能出现重复消费:如果 Consumer A 处理了 Offset 100-199 但只提交了 Offset 150(手动提交),Consumer B 会从 Offset 150 开始消费,导致 150-199 被重复处理。这是 At-least-once 语义下的正常行为。