主题
第 11 章 消费者组与 Rebalance:Kafka 客户端「最大的痛」
目标读者:写过 Consumer 但说不清「Rebalance 到底什么时候触发」「为什么我的消费在动态扩容时全停了 30 秒」「
group.instance.id是什么」「max.poll.interval.ms跟session.timeout.ms哪个先生效」的同学。学完你会:能画出 Consumer 加入 Group 的完整 4 步握手,能说清 Range/RoundRobin/Sticky/CooperativeSticky 的差别,知道 Eager Rebalance 为什么会「全停」、Cooperative 为什么能「增量」,能用
group.instance.id治住临时网络抖动引发的 Rebalance 风暴,遇到 Lag 持续上涨能按清单逐项排查。
0. 导读:消费者组是 Kafka 「最朴素却最容易踩坑」的设计
很多同学觉得 Consumer 写起来比 Producer 简单:subscribe、poll、处理消息、commit,循环往复。
但真到生产环境,几乎所有 Kafka 的「线上事故」里都有 Consumer 的影子:
- 「消费者重启一下,整组消费暂停了 30 秒」——Eager Rebalance;
- 「业务处理变慢了,消费者反复掉线,Lag 越涨越高」——
max.poll.interval.ms超时; - 「明明加了机器,但分区还是只分给原来那 3 个 Consumer」——Sticky 策略起作用;
- 「半夜网络抖动 5 秒,第二天早上发现 Lag 涨了好几个亿」——Rebalance 风暴;
- 「消费者关掉再起来,几条消息被消费两遍」——offset 还没 commit 就崩了。
Rebalance 是这一切混乱的中心。它是 Kafka 客户端协议里最复杂的部分——本质上是个分布式协议,由 Group Coordinator(Broker 端)+ Consumer(客户端)共同参与,需要协调每个 Consumer 的「我能消费这些分区」与「整组分区怎么分」两个状态。
这一章把这盘乱麻一根一根抽出来。
1. Group Coordinator:消费组的「调度员」
1.1 消费组(Consumer Group)是什么
消费组 = 一组
group.id相同的 Consumer,它们共享对一个 Topic 的消费工作。同一组里每个分区只会被一个 Consumer 消费。
举例:Topic orders 有 6 个分区,消费组 pay-service 里有 3 个 Consumer:
分区 0,1 → Consumer-A
分区 2,3 → Consumer-B
分区 4,5 → Consumer-C如果 pay-service 增加到 6 个 Consumer,每个就分到 1 个分区。如果增加到 8 个,会有 2 个 Consumer 空转(没分到分区)。这是为什么有句口诀:
消费者并行度不会超过分区数。
1.2 谁来分配?Group Coordinator
Kafka 不需要外部协调系统来管消费组。它把这件事交给一个特殊角色——Group Coordinator:
Group Coordinator = 处理这个 Group 的「调度员」,负责接收 Consumer 的加入请求、计算分区分配、维护心跳、提交 offset。
Group Coordinator 是某一台 Broker。怎么决定是哪一台?
coordinatorBrokerId = __consumer_offsets 主题中
hash("group.id") % numPartitions(__consumer_offsets)
所对应分区的 Leader 所在的 Broker也就是说:
- Kafka 内部有个 Topic
__consumer_offsets,默认 50 个分区; - 把消费组名 hash,映射到某个分区;
- 那个分区的 Leader Broker 就是这个消费组的 Group Coordinator。
生活类比:每家公司有个总机,你要找哪个部门,先按部门名首字母拨内线,电话转到对应的接线员。Coordinator 就是那个接线员,决定了哪个 Consumer 接哪条线。
1.3 Coordinator 也会切换
如果 Coordinator 所在 Broker 宕机,那个 __consumer_offsets 分区会经历 ISR 切换,新 Leader 上的那台 Broker 就成了新的 Coordinator。Consumer 会收到 NOT_COORDINATOR 错误,自动 FindCoordinator 重新定位,对业务透明。
📌 小心:__consumer_offsets 副本数默认 = offsets.topic.replication.factor(默认 3,但单 Broker 集群会被强制设为 1)。生产中一定要保证副本数 ≥ 3,否则 Coordinator 节点宕了会导致整组消费暂停。
2. 消费组协议四大请求:JoinGroup / SyncGroup / Heartbeat / LeaveGroup
整个消费组协议有 4 个核心请求。理解它们就能讲清楚「Consumer 加入 / 维持 / 退出 Group」的全过程。
2.1 时序图:一个 Consumer 加入 Group
2.2 四个请求各自做什么
| 请求 | 谁发 | 谁回 | 干什么 |
|---|---|---|---|
FindCoordinator | Consumer | 任意 Broker | 通过 hash(group.id) 找到 Coordinator 的 broker |
JoinGroup | Consumer | Group Coordinator | 「我要加入这个组」,注册自己支持的分配策略列表 |
SyncGroup | Consumer (Group Leader) | Group Coordinator | Group Leader 把计算好的分配方案告知 Coordinator 转发 |
Heartbeat | Consumer | Group Coordinator | 「我还活着」,Coordinator 也借此通知有 Rebalance 发生 |
LeaveGroup | Consumer | Group Coordinator | 优雅下线,触发 Rebalance(不发也行,但要等心跳超时) |
OffsetCommit | Consumer | Group Coordinator | 提交 offset 到 __consumer_offsets Topic |
OffsetFetch | Consumer | Group Coordinator | 启动时拉自己负责的分区上一次的 offset |
2.3 Group Leader vs Coordinator(容易混)
- Group Coordinator 是 Broker 端的角色——某台 Broker 兼任,负责接受请求、维护状态。
- Group Leader 是 Consumer 端的角色——同一个 Group 中第一个 join 进来的 Consumer 被指定为 Leader,它在客户端本地计算「谁分到哪些分区」,然后把方案通过 SyncGroup 上报给 Coordinator。
为什么把分配算法放在客户端? 这是 Kafka 一个聪明的设计:分配算法可以独立演进、用户可以自定义、Broker 端不需要理解协议细节。Broker 只是个透传方案的「邮局」。
3. 触发 Rebalance 的所有情形
记住一个总规则:消费组成员构成 / 订阅关系 / 分区数发生变化时,都会触发 Rebalance。具体:
- 新 Consumer 加入 Group(启动 / 扩容);
- Consumer 离开 Group(主动 LeaveGroup / 进程崩溃 / 心跳超时 / 处理超时);
- 订阅的 Topic 变化(
subscribe(["a", "b"])改成subscribe(["a", "b", "c"])); - 订阅的 Topic 分区数变化(管理员扩了分区);
- 正则订阅
subscribe(Pattern)时,匹配到了新创建的 Topic。
每次 Rebalance,Coordinator 把 generation_id +1。所有旧 generation 的请求都会被拒绝。
3.1 关键的两个超时
session.timeout.ms 默认 45000 Coordinator 多久没收到 Heartbeat 就认为 Consumer 死了
heartbeat.interval.ms 默认 3000 Consumer 多久发一次 Heartbeat
max.poll.interval.ms 默认 300000 (5 分钟) 两次 poll() 之间最长间隔最容易踩的坑是 max.poll.interval.ms:
- Consumer 内部有个后台心跳线程,定时给 Coordinator 发 Heartbeat(即使主线程在处理消息);
- 但主线程必须按时回到
poll(),否则 Coordinator 会通过另一个机制(poll 超时)认为这个 Consumer 处理太慢,把它踢出组; - 也就是说,session.timeout.ms 管「网络/进程是否活着」,max.poll.interval.ms 管「业务线程是否活着」。
典型事故剧本:你的消费逻辑某次需要 7 分钟(比如调外部接口超时)。后台心跳一直在发,Coordinator 不会因为 session.timeout 把你踢;但 5 分钟过后,poll 超时了——Coordinator 把你踢出组,触发 Rebalance;你处理完那条消息回来调 poll(),Broker 告诉你「你已经不在 Group 里了」,导致:
- 整组 Rebalance 30 秒(其它 Consumer 也跟着停);
- 你处理完的那条消息可能要被另一个 Consumer 重消费(如果你还没 commit);
- 你的 Consumer 重新 join 又触发一次 Rebalance;
- 业务逻辑慢的根因没解决,下一批消息又会超时——死循环 Rebalance 风暴。
修复思路:调大 max.poll.interval.ms 是治标,真正的根治是把每条消息处理时间砍下来——拆成异步、加并发、降批量。
4. 四大分区分配策略对比
Kafka 内置 4 种分配策略(partition.assignment.strategy):
RangeAssignor(默认,单 Topic 范围分)RoundRobinAssignor(轮询)StickyAssignor(粘性,2.4 之前)CooperativeStickyAssignor(合作式粘性,2.4+ 推荐)
下面用统一例子演示:3 个 Consumer(C1, C2, C3)订阅 1 个 Topic(6 个分区 P0~P5)。
4.1 Range:每个 Topic 内按范围切
按 Consumer 字典序、Topic 内分区序,分段切给每个 Consumer:
Topic orders 有 6 分区, 3 Consumer
6 / 3 = 2,每人 2 个
C1 ← P0, P1
C2 ← P2, P3
C3 ← P4, P5问题:当订阅多个 Topic 时,每个 Topic 都按这个规则单独分,可能导致严重不均:
订阅 3 个 Topic,每个都 6 分区:
C1 ← topicA[P0,P1] + topicB[P0,P1] + topicC[P0,P1] = 6 分区
C2 ← topicA[P2,P3] + topicB[P2,P3] + topicC[P2,P3] = 6 分区
C3 ← topicA[P4,P5] + topicB[P4,P5] + topicC[P4,P5] = 6 分区看起来均衡?但如果 Topic 有 7 分区:
C1 ← topicA[P0,P1,P2] + topicB[P0,P1,P2] + topicC[P0,P1,P2] = 9 分区
C2 ← topicA[P3,P4] + topicB[P3,P4] + topicC[P3,P4] = 6 分区
C3 ← topicA[P5,P6] + topicB[P5,P6] + topicC[P5,P6] = 6 分区C1 比别人多 50%。Range 默认策略至今未变,但生产环境很少用。
4.2 RoundRobin:所有分区按轮询
把订阅的所有 (topic, partition) 对展开排序,再轮流发给 Consumer:
订阅 1 个 Topic 6 分区:
P0→C1, P1→C2, P2→C3, P3→C1, P4→C2, P5→C3
等价于:
C1 ← P0, P3
C2 ← P1, P4
C3 ← P2, P5优点:跨 Topic 也均匀。 缺点:所有 Consumer 必须订阅完全相同的 Topic 列表,否则分配会出乱。
4.3 Sticky:尽量保留旧分配
StickyAssignor(KIP-54)做了两件事:
- 均匀(与 RoundRobin 一样均匀,不一定一致);
- 粘性:Rebalance 时尽量保留每个 Consumer 原来的分区分配,减少不必要的迁移。
对比例子(C2 离组):
原分配: C1=[P0,P1] C2=[P2,P3] C3=[P4,P5]
Range / RoundRobin 在 Rebalance 时可能重洗:
C1=[P0,P1,P2] C3=[P3,P4,P5]
→ C1 多接管了 P2,C3 拿到了陌生的 P3
Sticky 重排:
C1=[P0,P1, P2] C3=[P4,P5, P3]
→ 保持了 P0,P1,P4,P5 不动,只接管 C2 留下的 P2,P3「保持不动」的好处:本地缓存(去重表、聚合状态)不必清空重建。
4.4 CooperativeSticky:增量 Rebalance(2.4+ 推荐)
CooperativeStickyAssignor(KIP-429)把「Sticky」从「分配结果」推进到「Rebalance 过程」本身。
Eager 协议的痛
传统的 Eager Rebalance 是「Stop the World」式:
1. Coordinator 发起 Rebalance
2. 所有 Consumer 撤销(revoke)当前所有分区
3. 所有 Consumer 重新 JoinGroup + SyncGroup
4. 拿到新分配,重新开始消费第 2 步开始到第 4 步结束的整个时间段——整组 Consumer 全部停止消费。即使分配结果与原来 90% 一样,也要全停。
Cooperative 协议的「增量」
Cooperative Rebalance 把这个过程拆成两轮:
第一轮 Rebalance:
1. Coordinator 发起 Rebalance
2. Consumer 仍然继续消费当前分区(不撤销)
3. JoinGroup + SyncGroup → 算出新分配方案
4. Consumer 比较新旧分配:
- 仍属于自己的分区:继续消费,不动
- 自己即将失去的分区:调用 onPartitionsRevoked() 撤销
5. 撤销完成后再触发第二轮 Rebalance
第二轮 Rebalance:
6. 重新 JoinGroup + SyncGroup
7. Consumer 拿到新分到的分区,调用 onPartitionsAssigned()直观体感:稳态 Rebalance 时,90% 的 Consumer 完全不停,只有少数 Consumer 经历短暂的「revoke 旧 → assign 新」。整体停顿时间从「秒级」降到「毫秒级」。
启用方式
python
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'pay-service',
'partition.assignment.strategy': 'cooperative-sticky', # 关键!
...
})⚠️ 滚动升级注意:从 Eager 切到 Cooperative 不能一步到位,要走过渡。先把策略改成
[cooperative-sticky, range]多策略列表(让新老 Consumer 共存时还能选 range),全部 Consumer 升级完再去掉 range,单纯保留 cooperative-sticky。
4.5 四种策略对比表
| 策略 | 是否均衡 | 是否粘性 | 是否增量 | 推荐场景 |
|---|---|---|---|---|
| Range(默认) | 跨 Topic 不均匀 | ❌ | ❌ | 单 Topic 简单场景 |
| RoundRobin | 均匀(要求订阅一致) | ❌ | ❌ | 多 Topic 简单场景 |
| Sticky | 均匀 | ✅ | ❌ | 2.4 前的 best practice |
| CooperativeSticky | 均匀 | ✅ | ✅ | 2.4+ 强烈推荐 |
5. Static Membership:治住「假性 Rebalance」
5.1 痛点:临时网络抖动 = Rebalance 风暴
经典场景:
- 消费者 C 跟 Coordinator 之间网络抖了 6 秒,Coordinator 没收到心跳,认为 C 死了;
- 触发 Rebalance,剩余 Consumer 接管 C 的分区;
- 等 C 网络恢复,发现自己已经不在组里了,重新 JoinGroup → 又触发一次 Rebalance;
- 整组消费在这两次 Rebalance 期间停顿数十秒。
根因:Kafka 默认用「member.id」标识 Consumer,每次 join 都生成新的随机 ID。Coordinator 没法识别「这个 Consumer 是不是刚才那个」,只能按「新成员加入」处理。
5.2 解决:group.instance.id(Kafka 2.3+)
KIP-345 引入 Static Membership:给每个 Consumer 配一个人工指定的稳定 ID:
python
consumer = Consumer({
'bootstrap.servers': '...',
'group.id': 'pay-service',
'group.instance.id': 'consumer-pod-3', # ← 关键
'session.timeout.ms': 60000, # 配套调大,给重连留空间
...
})这样,Coordinator 在 session.timeout.ms 内:
- 看不到
consumer-pod-3心跳?不立刻 Rebalance,先等等; - 在超时之前
consumer-pod-3又出现并 join 了?直接把它原来的分区还给它,不发生 Rebalance; - 真的超时了?才触发 Rebalance。
5.3 配套配置
group.instance.id:每个 Consumer 实例全局唯一(K8s 里通常用 Pod name + 序号);session.timeout.ms:调大到 30~60 秒,给临时抖动留窗口;- 下线时主动调用
consumer.close(autoCommit=true)优雅退出,否则要等 session.timeout 才触发 Rebalance。
📌 副作用:调大 session.timeout 意味着「真死的 Consumer 也要更久才被发现」。要根据业务能容忍的最大停顿时间权衡。
6. Rebalance 监听器:资源清理范式
python
from confluent_kafka import Consumer, ConsumerRebalanceListener
class MyListener:
def on_assign(self, consumer, partitions):
# Rebalance 后被分配到这些分区时调用
# 适合:初始化每个分区的本地缓存 / 加载位点 / 启动监控
for tp in partitions:
print(f"[on_assign] 接管 {tp.topic}-{tp.partition}")
def on_revoke(self, consumer, partitions):
# Rebalance 前撤销自己持有的分区时调用
# 适合:commit 当前 offset / 清理本地缓存 / 关闭外部连接
try:
consumer.commit(asynchronous=False)
except Exception as e:
print(f"[on_revoke] commit 失败:{e}")
for tp in partitions:
print(f"[on_revoke] 撤销 {tp.topic}-{tp.partition}")
def on_lost(self, consumer, partitions):
# 仅 cooperative-sticky 协议:被强制 revoke(自己不再有所有权)
# 不应该 commit(offset 已经无效)
for tp in partitions:
print(f"[on_lost] 失去 {tp.topic}-{tp.partition}")
consumer = Consumer({...})
consumer.subscribe(['orders'],
on_assign=MyListener().on_assign,
on_revoke=MyListener().on_revoke,
on_lost=MyListener().on_lost)6.1 黄金法则
- on_revoke 一定要 commit(同步 commit,不能 fire-and-forget),否则同一条消息会被新接手的 Consumer 重消费;
- on_revoke 是阻塞的:你在里面做的任何事都会延长 Rebalance 时间;不要在里面调外部接口、做大量 IO;
- on_lost 不要 commit:你已经没有这些分区的所有权,commit 会被拒;
- on_assign 里不要重置 offset(除非你明确知道要重新消费);默认 Coordinator 会从
__consumer_offsets给你最近一次 commit 的位置。
7. 消费 Lag 排查实战
7.1 用 kafka-consumer-groups.sh --describe 查 Lag
bash
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group pay-service输出示例:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
pay-service orders 0 12345 12345 0 consumer-pay-service-1-uuid /10.0.0.1 consumer-1
pay-service orders 1 54300 58000 3700 consumer-pay-service-2-uuid /10.0.0.2 consumer-2
pay-service orders 2 0 1000 1000 - - -读法:
- CURRENT-OFFSET:消费组在该分区已 commit 的 offset;
- LOG-END-OFFSET:该分区最新写入的 offset;
- LAG = LOG-END-OFFSET - CURRENT-OFFSET:未消费的消息数;
- CONSUMER-ID = "-":表示这个分区当前没人消费(要么没 Consumer 启动,要么订阅没生效)。
7.2 Lag 持续上涨的根因清单(按出现频率排)
1. 消费速度跟不上生产速度(最常见)
排查:看 Producer QPS 与 Consumer QPS 对比;调 max.poll.records(一次 poll 拉多少条)、增加 Consumer 实例(注意不超过分区数)。
2. 业务处理慢,导致 max.poll.interval.ms 超时 → Rebalance 风暴
排查:看 server.log 是否高频出现 "Member ... has failed, removing it from the group";看消费侧日志是否反复 "Revoking..."。
修复:把每条消息处理时间砍下来(加并发、异步化),别把消息处理逻辑跟 IO 串行。
3. 单 Consumer 实例处理能力打满
排查:看 Consumer 进程 CPU / IO / GC;可能是单线程消费 + 重 CPU 操作。
修复:每个 Consumer 内部用线程池处理消息,但要注意 offset 不能在异步任务里直接 commit(容易乱序),常见模式是「主线程 poll + 提交、子线程并发处理」。
4. 部分分区积压(Key 热点)
排查:上面的命令里看 LAG 列,是不是只有 1~2 个分区在涨?很可能 Producer 端某个 key 的消息量远超其它,导致该分区积压。
修复:调整 Producer 的 partitioner,把热 key 拆成多个 sub-key;或者业务侧用更细粒度的 key。
5. Consumer 实际未启动(CONSUMER-ID 为 "-")
排查:上面命令的 CONSUMER-ID 那列是 - 说明没有 Consumer 在订阅;可能是部署没起来,或是订阅写错了 Topic。
6. 消费组在 Rebalance 中
排查:--state 选项查看消费组当前状态:
bash
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group pay-service --state输出有 STATE 列:Stable / PreparingRebalance / CompletingRebalance / Empty / Dead。如果长期不是 Stable,要查为什么频繁 Rebalance。
7. Coordinator 那台 Broker 故障
排查:__consumer_offsets 主题对应分区的 Leader 在哪个 Broker?那台是不是 GC 严重 / 磁盘满 / 进程挂了?
8. 动态加分区但消费组没感知
Kafka 默认 5 分钟刷新一次元数据(metadata.max.age.ms)。如果 Topic 加了分区但 Consumer 没刷新,新分区不会被消费。可以重启 Consumer 强刷。
9. 反向:Lag 显示极大但实际无消息(虚假 Lag)
某些场景(比如长期不消费的旧消费组)会显示 Lag = 几百亿——其实是 LOG-END-OFFSET 在涨而消费组从未提交过,offset 默认从 latest 算起就显示「全部未消费」。重置一次 offset 即可:
bash
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group pay-service --reset-offsets --to-latest \
--topic orders --execute8. 配置速查表
8.1 客户端必配
| 配置 | 默认 | 推荐 | 说明 |
|---|---|---|---|
group.id | - | 业务名 | 消费组标识 |
partition.assignment.strategy | range | cooperative-sticky | 强烈推荐 |
enable.auto.commit | true | false | 手动 commit 才能保证不丢 |
auto.offset.reset | latest | earliest(首次/重置) | 没有 commit 时从哪开始 |
session.timeout.ms | 45000 | 30000~60000 | 配 static membership 时调大 |
heartbeat.interval.ms | 3000 | session.timeout / 3 | Heartbeat 频率 |
max.poll.interval.ms | 300000 | 看业务最长处理时间 | 业务慢就调大或拆消费 |
max.poll.records | 500 | 视消息大小 | 一次 poll 最多拿多少条 |
fetch.min.bytes / max.wait.ms | 1 / 500 | 高吞吐场景调大 | Broker 端聚合等待 |
group.instance.id | - | 用 Pod 名 | 静态成员,治抖 |
8.2 服务端相关
| 配置 | 默认 | 说明 |
|---|---|---|
offsets.topic.replication.factor | 3 | __consumer_offsets 副本数(务必 ≥3) |
offsets.topic.num.partitions | 50 | 分区数(决定能承载多少消费组并发) |
group.initial.rebalance.delay.ms | 3000 | 初始 Rebalance 等待,给批量加入留时间 |
9. 与其他 MQ 的「消费组」对比
┌──────────────┬──────────────────────────┬───────────────────────────┐
│ 系统 │ 消费组模型 │ Rebalance / 协调 │
├──────────────┼──────────────────────────┼───────────────────────────┤
│ Kafka │ Group 内分区独占 │ Coordinator + JoinGroup │
│ │ 「拉」模式 + 手动 offset │ 4 种策略,CooperativeSticky│
│ RabbitMQ │ Queue 多 Consumer 竞争 │ 无(推模式 + ack) │
│ │ 「推」模式 + ACK │ │
│ RocketMQ │ 集群 / 广播模式 │ Rebalance(5.0 改进) │
│ │ 队列 vs 分区粒度 │ │
│ Pulsar │ 4 种订阅模式 │ Broker 端 push 路由 │
│ │ Exclusive/Shared/... │ 不需要客户端 Rebalance │
│ NATS JetStream│ Pull/Push Consumer │ Server 端管理 offset │
└──────────────┴──────────────────────────┴───────────────────────────┘Kafka 的独占模型在有序性和回放能力上无可替代;代价就是 Rebalance 这套复杂协议。RabbitMQ 之类的推模式没有 Rebalance 概念,但代价是「消息一旦推给某个 Consumer 处理失败就要 requeue,更难保序」。
10. 小结:Rebalance 治理的 5 条军规
把整章压缩成 5 条:
- 首选 cooperative-sticky 策略:从 Kafka 2.4 起就有,能省 80% 的「全停时间」;
- 配
group.instance.id+ 调大session.timeout.ms:治住临时抖动引发的「假性 Rebalance」; - 业务处理时间不要超过
max.poll.interval.ms:宁可拆细每条消息的处理,也别让 poll 超时; - 手动 commit + on_revoke 同步 commit:防止重复消费;
- 消费 Lag 接监控告警:盯住 LAG 与 STATE 两个指标,频繁 Rebalance 比 Lag 上涨更可怕。
11. 配套代码与演示
11_consumer_group/code/cooperative_consumer.py:partition.assignment.strategy=cooperative-sticky的最简 Consumer 范本;11_consumer_group/code/static_membership.py:演示group.instance.id在「关停 Consumer 5 秒再起来」时不触发 Rebalance;11_consumer_group/code/lag_monitor.py:用 AdminClient 周期打印每个分区的 Lag;11_consumer_group/code/rebalance_listener.py:自定义on_assign / on_revoke / on_lost,演示资源清理;11_consumer_group/demo.html:可视化 4 种分区分配策略对比、Eager 与 Cooperative Rebalance 时序、Static Membership 开关效果。
12. 面试高频题
Q1:Group Coordinator 是怎么确定的?切换会发生什么?
考察点:消费组协议的存储与定位。
标准答案:
- Kafka 集群里有个内部 Topic
__consumer_offsets(默认 50 分区); - 给定消费组名,
hash(group.id) % 50算出对应分区; - 那个分区的 Leader Broker 就是这个消费组的 Coordinator;
- 当 Coordinator 节点宕机:
__consumer_offsets那个分区会经历 ISR 切换,新 Leader 上线;- Consumer 发请求到旧 Coordinator 收到
NOT_COORDINATOR错误; - Consumer 重新调用
FindCoordinator定位到新 Coordinator; - Consumer 重连后发 JoinGroup,触发一次 Rebalance;
- 所以
__consumer_offsets的副本数(offsets.topic.replication.factor)一定要 ≥ 3。
Q2:详细描述一个 Consumer 加入 Group 的全过程
考察点:协议层的细节理解。
标准答案:四步:
- FindCoordinator:Consumer 拿
group.id询问任意 Broker,返回 Coordinator 所在 Broker; - JoinGroup:Consumer 向 Coordinator 发 JoinGroup(带支持的策略列表),Coordinator 等够
group.initial.rebalance.delay.ms后把全员状态返回;返回时指定第一个 join 的 Consumer 为 Group Leader; - SyncGroup:Group Leader 在客户端本地按选定策略计算分配方案,把方案随 SyncGroup 上报;其它 Consumer 也发 SyncGroup(不带方案),Coordinator 收齐后把方案分发给每个 Consumer;
- Heartbeat 循环:Consumer 之后每
heartbeat.interval.ms给 Coordinator 发心跳,正常返回 ok;如果有 Rebalance 在进行,返回 REBALANCE_IN_PROGRESS,Consumer 知道要重新 JoinGroup。
加分项:补充「分配算法在客户端而不是 Broker,是因为可独立演进、用户可自定义」。
Q3:CooperativeSticky 比 Eager Rebalance 强在哪?
考察点:对 Kafka 2.4+ 新协议的理解。
标准答案:
- Eager Rebalance:所有 Consumer 在 Rebalance 开始时全部撤销所有分区,等新分配下来才开始消费——整组停顿,即使 99% 分配结果与原来一样也要全停。
- Cooperative Rebalance:分两轮,
- 第一轮:Consumer 仍正常消费,Coordinator 算出新分配;Consumer 比较新旧,只撤销那些自己即将失去的分区;
- 第二轮:撤销完成后再走一次 JoinGroup,把空出的分区分给新 Consumer;
- 效果:稳态 Rebalance 时绝大多数 Consumer完全不停,只有极少数经历短暂 revoke→assign。整体停顿从「秒级」降到「毫秒级」。
- 启用方式:
partition.assignment.strategy=cooperative-sticky(2.4+)。 - 滚动升级注意:要先用
[cooperative-sticky, range]多策略列表过渡,全员升完再去掉 range。
Q4:session.timeout.ms 和 max.poll.interval.ms 有什么区别?哪个先生效?
考察点:Consumer 心跳与处理超时的边界。
标准答案:
session.timeout.ms(默认 45s):管「Consumer 进程或网络是否还活着」。Consumer 有一个后台心跳线程定时发 Heartbeat,Coordinator 收不到心跳超过 session.timeout 就认为 Consumer 死了。max.poll.interval.ms(默认 5min):管「Consumer 业务线程是否还活着」。Consumer 主线程如果两次poll()之间间隔超过这个值,Coordinator 也会把它踢出组——即使心跳还在正常发。- 谁先生效:取决于谁先超时。但实战中绝大多数生产事故都是 max.poll.interval.ms 超时——因为业务处理变慢的概率,比网络/进程突然挂掉高得多。
- 治本办法:把每条消息处理时间砍下来(异步化 / 加并发 / 拆消息),不要单纯调大 max.poll.interval.ms。
Q5:group.instance.id(Static Membership)解决什么问题?
考察点:对 KIP-345 新特性的认识。
标准答案:
- 传统问题:Consumer 每次启动随机生成
member.id。如果 Consumer 短暂网络抖动让 Coordinator 误以为它死了,会触发 Rebalance;等 Consumer 恢复时被当成「新成员」加入,又触发一次 Rebalance——抖动 5 秒可能引发 30 秒的双 Rebalance。 - 解决方案:Kafka 2.3+ 引入
group.instance.id:用户给每个 Consumer 实例配一个人工指定的稳定 ID(K8s 里通常是 Pod name)。 - 效果:在
session.timeout.ms内,Coordinator 看不到这个 instance.id 的心跳不会立刻触发 Rebalance;如果在超时之前 Consumer 又出现并 join,Coordinator 直接把原来的分区还给它,整组无 Rebalance。 - 代价:调大 session.timeout 意味着「真正的故障检测」变慢了。要根据业务能容忍的最大停顿时间权衡。
Q6:消费 Lag 持续上涨怎么排查?
考察点:实战能力。
标准答案(按概率从高到低):
- 生产>消费:看 Producer QPS 和 Consumer QPS 对比,扩 Consumer 实例、调大 max.poll.records;
- 业务处理太慢导致 max.poll.interval.ms 超时:看 server.log 是否高频 Rebalance;改进消息处理逻辑;
- 单 Consumer 实例打满:看 CPU/GC/IO;
- Key 热点导致部分分区积压:用
--describe看每分区 LAG; - CONSUMER-ID 是
-:根本没人在消费; - 消费组在 Rebalance:看
--state,长期非 Stable 要查原因; - Coordinator 节点故障:检查
__consumer_offsets对应分区的 Leader Broker; - Topic 加了分区但 Consumer 没感知:等 metadata.max.age.ms 或重启 Consumer。
排查工具:kafka-consumer-groups.sh --describe、JMX records-lag-max、Kafka UI。
Q7:自动 commit 和手动 commit 各有什么风险?
考察点:消息语义。
标准答案:
- 自动 commit(
enable.auto.commit=true,默认每 5s 提交一次):- 风险 1(重复消费):处理了消息但还没到自动提交的时刻就崩溃,重启会从上一次自动 commit 的位置开始 → 重复消费;
- 风险 2(消息丢失):自动 commit 是「按上次 poll 拿到的最大 offset」提交,poll 之后还在处理,commit 已经过去了——如果业务处理失败但 commit 已发,那条消息就丢了;
- 手动 commit:
- 同步 commit (
commitSync):阻塞,可靠但慢; - 异步 commit (
commitAsync):不阻塞,但出错没回调就忘了; - 推荐组合:日常
commitAsync+ 关闭/Rebalance 时commitSync+ on_revoke 里同步 commit;
- 同步 commit (
- 真正的 Exactly Once 需要 Producer + Consumer 配合:用 Kafka 事务(
sendOffsetsToTransaction)把「消费 offset 的提交」也纳入同一个事务(详见第 13 章)。
Q8:为什么消费组里 Consumer 实例数不能超过分区数?
考察点:消费模型的本质。
标准答案:
- Kafka 的消费模型规定同一组内每个分区只能被一个 Consumer 消费;
- 所以 N 个分区最多支持 N 个 Consumer 同时工作;
- 多出来的 Consumer 会被分配到「空集合」——它们建立 TCP 连接、参与 JoinGroup / Heartbeat、消耗内存与网络,但不消费任何消息;
- 扩容时正确做法是先扩分区再扩 Consumer;
- 但分区数也不是越多越好:分区是文件句柄、副本同步线程、Controller 元数据的开销单位;通常单 Topic 几百分区就足够,超过几千要谨慎评估。
加分项:补充「Pulsar 用 Shared Subscription 不受这个限制——同一个分区可以被多个 Consumer 并行消费,代价是顺序保证退化」。
至此,9 / 10 / 11 三章打通了 Kafka 的「副本机制」「元数据管理」「消费协议」三个核心模块。下一章我们将正式进入消息语义与 Exactly Once 的世界。
13. 拓展:常见误解与「冷知识」
13.1 误解 ①:「Rebalance 时 Producer 也要停」
不会。Producer 与 Consumer 完全独立。Rebalance 只发生在 Consumer Group 内部,对 Producer 透明。Producer 在 Rebalance 期间继续往分区写入,分区里的消息正常累积;Rebalance 一结束,新接管该分区的 Consumer 会从上次 commit 的 offset 接着往下消费。
13.2 误解 ②:「Consumer 数 = 分区数 时性能最高」
不一定。Consumer 实例数等于分区数意味着每个 Consumer 一对一负责一个分区,这是「每个分区都有人消费」的最少 Consumer 数。但单个 Consumer 是否打满,取决于消息处理逻辑:
- 处理逻辑很轻(< 1ms):1 个 Consumer 一个分区也够;
- 处理逻辑重(IO / 调用外部接口):单个 Consumer 配多线程处理才能打满分区吞吐。
「Consumer 数 = 分区数」只是「单个 Consumer 内不再开多线程」的边界条件,不是性能最高的边界。
13.3 冷知识:消费组协议的 Generation
每次 Rebalance 完成,Coordinator 会把 Group 的 generation_id +1。所有协议请求都要带上自己的 generation:
- Consumer 提交 offset 时,generation 与 Coordinator 不一致 → 提交失败;
- Consumer 发 Heartbeat 时,generation 不一致 → 收到 ILLEGAL_GENERATION 错误,必须重新 JoinGroup;
生活类比:generation 就像「值日表的版本号」。新版本一发布,旧版本的所有请求都作废,必须用新版本重排。这是为什么频繁 Rebalance 不仅停顿严重,还会导致大量 offset 提交失败、消息重复消费。
13.4 冷知识:__consumer_offsets 是 Compact Topic
__consumer_offsets 用的是 Log Compaction 保留策略(不是普通的按时间删除):
- 同一个 (group, topic, partition) 的 offset 提交,新值会覆盖旧值;
- Compaction 周期性执行,旧的 offset 记录被清理,只留每个 key 的最新一条;
- 所以即便消费组运行了几年,
__consumer_offsets也不会无限膨胀。
详见第 14 章 Log Compaction。
13.5 冷知识:subscribe(Pattern) 的副作用
用正则订阅 Topic(如 subscribe(Pattern.compile("logs\\..+")))时,Consumer 会周期性地刷新 metadata(metadata.max.age.ms)来发现新匹配的 Topic。每发现一个新 Topic,整个 Group 触发一次 Rebalance。
如果你的命名规则容易匹配到大量动态创建的 Topic(比如 CDC 场景下每张新表都创建一个 Topic),会导致频繁 Rebalance。生产中能列举就别用正则。
14. 速记口诀
最后送几句顺口溜帮你记关键点:
- 「四请求八步骤」 —— FindCoordinator / JoinGroup / SyncGroup / Heartbeat 四请求,配合 OffsetCommit/Fetch/LeaveGroup 八步骤;
- 「Coordinator 在 Broker,Leader 在客户端」 —— 别搞混 Group Coordinator 和 Group Leader;
- 「分区数是上限」 —— Consumer 实例数永远别超过分区数;
- 「session 管命,poll 管业」 —— session.timeout 检测进程死活,max.poll.interval 检测业务死活;
- 「优先 Cooperative」 —— 2.4+ 都该用 cooperative-sticky;
- 「Static Membership 治抖动」 —— Pod 网络飘忽就配 group.instance.id;
- 「revoke 必 commit」 —— Rebalance 监听器 on_revoke 一定要同步 commit;
- 「Lag 高不可怕,频繁 Rebalance 才可怕」 —— 监控指标排序记牢。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
#!/usr/bin/env python3
"""
cooperative_consumer.py — 演示 cooperative-sticky 分配策略
观察重点:
1) 启动一个 Consumer 时打印 on_assign 的分区列表;
2) 启动第二个 Consumer(同 group.id)时,第一个 Consumer 只会 revoke 一部分分区,
不会经历「全部撤销 → 重新分配」的 Stop-the-World;
3) 关闭其中一个 Consumer,剩下的仅增量接管,不全停。
用法:
# 终端 1
python cooperative_consumer.py --consumer-id c1
# 终端 2(5 秒后启动)
python cooperative_consumer.py --consumer-id c2
# 终端 3
python cooperative_consumer.py --consumer-id c3
前置:
Topic 至少 3 分区
kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic learn.11.coop --partitions 6 --replication-factor 3 \
--if-not-exists
生产端可以用 kafka-console-producer.sh 持续写入:
kafka-console-producer.sh --bootstrap-server localhost:9092 --topic learn.11.coop
"""
from __future__ import annotations
import argparse
import signal
import sys
import time
from datetime import datetime
from confluent_kafka import Consumer, TopicPartition
def parse_args():
p = argparse.ArgumentParser()
p.add_argument("--bootstrap", default="localhost:9092")
p.add_argument("--topic", default="learn.11.coop")
p.add_argument("--group", default="coop-demo-group")
p.add_argument("--consumer-id", default="c1")
return p.parse_args()
def fmt_tps(tps):
return ", ".join(f"{tp.topic}-{tp.partition}" for tp in tps) or "<empty>"
def main():
args = parse_args()
cid = args.consumer_id
def log(msg, color="\033[36m"):
print(f"{color}[{datetime.now().strftime('%H:%M:%S')}][{cid}] {msg}\033[0m",
flush=True)
def on_assign(consumer, partitions):
# cooperative 模式下:partitions 是「新增」的分区,不是「全集」
log(f"on_assign → +{fmt_tps(partitions)}", "\033[32m")
consumer.incremental_assign(partitions)
def on_revoke(consumer, partitions):
# cooperative 模式下:partitions 是「失去」的分区,必须用 incremental_unassign
log(f"on_revoke → -{fmt_tps(partitions)}", "\033[33m")
try:
# 同步 commit 当前 offset,防止重复消费
consumer.commit(asynchronous=False)
except Exception as e:
log(f"commit failed: {e}", "\033[31m")
consumer.incremental_unassign(partitions)
def on_lost(consumer, partitions):
# 仅 cooperative 才会调用:所有权已经被 Coordinator 强制收回,不能 commit
log(f"on_lost → -{fmt_tps(partitions)} (forced)", "\033[31m")
consumer.incremental_unassign(partitions)
consumer = Consumer({
"bootstrap.servers": args.bootstrap,
"group.id": args.group,
"client.id": cid,
"partition.assignment.strategy": "cooperative-sticky", # ★关键★
"enable.auto.commit": False,
"auto.offset.reset": "latest",
"session.timeout.ms": 45000,
"heartbeat.interval.ms": 3000,
"max.poll.interval.ms": 300000,
})
consumer.subscribe([args.topic],
on_assign=on_assign,
on_revoke=on_revoke,
on_lost=on_lost)
log(f"started, subscribing to {args.topic} (cooperative-sticky)")
# 优雅退出
stop = {"flag": False}
def _sig(*_):
log("got signal, draining…", "\033[35m")
stop["flag"] = True
signal.signal(signal.SIGINT, _sig)
signal.signal(signal.SIGTERM, _sig)
last_print = time.time()
msg_count = 0
try:
while not stop["flag"]:
msg = consumer.poll(timeout=1.0)
if msg is None:
if time.time() - last_print > 5:
last_print = time.time()
held = consumer.assignment()
log(f"heartbeat OK | holding: {fmt_tps(held)} | "
f"msgs since start: {msg_count}", "\033[34m")
continue
if msg.error():
log(f"err: {msg.error()}", "\033[31m")
continue
msg_count += 1
if msg_count % 20 == 0 or msg_count <= 5:
log(f"poll {msg.topic()}-{msg.partition()}@{msg.offset()} "
f"key={msg.key()!r} val={msg.value()!r}", "\033[36m")
# 模拟业务处理时间(很短)
# 业务处理后异步提交(实际生产推荐定期 commit + 关停时同步 commit)
consumer.commit(message=msg, asynchronous=True)
finally:
log("closing…")
consumer.close()
log(f"closed. total msgs = {msg_count}")
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
sys.exit(0)python
#!/usr/bin/env python3
"""
lag_monitor.py — 周期性打印某个消费组的每分区 Lag
实现思路:
1) 用 AdminClient.list_consumer_group_offsets() 拿消费组当前提交的 offset;
2) 用 AdminClient.list_offsets()(confluent-kafka 2.0+)拿每个分区的 LATEST offset;
3) Lag = LATEST - committed;
4) 同时打印消费组当前 STATE / 每个 member / 它们绑定的 client.id;
5) 周期性刷新(每 N 秒一次)。
用法:
python lag_monitor.py --group pay-service
python lag_monitor.py --group pay-service --interval 5
优势 vs kafka-consumer-groups.sh:
* 可以集成到自己的告警 / Prometheus exporter;
* 输出可以是 JSON 方便后续处理(--json);
* 可以同时监控多个组(重复传 --group)。
"""
from __future__ import annotations
import argparse
import json
import sys
import time
from datetime import datetime
from confluent_kafka import (
OffsetSpec,
TopicPartition,
)
from confluent_kafka.admin import (
AdminClient,
ConsumerGroupState,
)
def parse_args():
p = argparse.ArgumentParser()
p.add_argument("--bootstrap", default="localhost:9092")
p.add_argument("--group", action="append", required=True,
help="要监控的消费组(可重复传多个)")
p.add_argument("--interval", type=float, default=3.0)
p.add_argument("--json", action="store_true", help="JSON 输出便于 grep")
p.add_argument("--once", action="store_true")
return p.parse_args()
def fetch_state(admin: AdminClient, group: str):
"""拉取消费组的状态、members、各分区 committed offset"""
fut = admin.describe_consumer_groups([group])[group]
desc = fut.result(timeout=15)
state = desc.state
members = []
for m in desc.members:
ass_parts = []
if m.assignment:
for tp in m.assignment.topic_partitions:
ass_parts.append(f"{tp.topic}-{tp.partition}")
members.append({
"member_id": m.member_id[:30],
"client_id": m.client_id,
"host": m.host,
"instance_id": m.group_instance_id or "-",
"assignments": ass_parts,
})
# Committed offsets
fut = admin.list_consumer_group_offsets([
# 不指定 partitions 表示拉全部
# confluent-kafka API: ListConsumerGroupOffsetsRequest
])
# 这里用 list_consumer_group_offsets 的简化签名
# 不同版本 API 略有差异,下面用 describe 的方式做兼容
committed = {}
try:
from confluent_kafka.admin import ConsumerGroupTopicPartitions
req = ConsumerGroupTopicPartitions(group)
fut = admin.list_consumer_group_offsets([req])[group]
res = fut.result(timeout=15)
for tp in res.topic_partitions:
committed[(tp.topic, tp.partition)] = tp.offset
except Exception as e:
print(f"[warn] list_consumer_group_offsets 不可用,跳过 committed: {e}",
file=sys.stderr)
return state, members, committed
def fetch_latest_offsets(admin: AdminClient, tps):
if not tps:
return {}
spec = {tp: OffsetSpec.latest() for tp in tps}
fut = admin.list_offsets(spec)
out = {}
for tp, fu in fut.items():
try:
res = fu.result(timeout=15)
out[(tp.topic, tp.partition)] = res.offset
except Exception as e:
out[(tp.topic, tp.partition)] = -1
return out
def render(group, state, members, committed, latest, args):
"""打印一个组的状态"""
if args.json:
rows = []
for (t, p), off in sorted(committed.items()):
le = latest.get((t, p), -1)
lag = max(0, le - off) if (off >= 0 and le >= 0) else None
rows.append({"topic": t, "partition": p, "committed": off,
"log_end": le, "lag": lag})
print(json.dumps({"ts": datetime.now().isoformat(), "group": group,
"state": str(state), "members": members,
"lags": rows}, ensure_ascii=False))
return
print(f"\n=== group={group} state={state} members={len(members)} "
f"@{datetime.now().strftime('%H:%M:%S')} ===")
if members:
print(f" {'member':<32}{'client':<20}{'host':<20}{'instance':<14}{'assigned'}")
for m in members:
print(f" {m['member_id']:<32}{m['client_id']:<20}{m['host']:<20}"
f"{m['instance_id']:<14}{','.join(m['assignments']) or '<empty>'}")
else:
print(" (no active members)")
if not committed:
print(" (no committed offsets — 该组从未提交过)")
return
print(f"\n {'Topic':<26}{'Part':<6}{'Committed':<14}{'Log-End':<14}"
f"{'Lag':<10}{'Status'}")
print(" " + "-" * 90)
total_lag = 0
for (t, p), off in sorted(committed.items()):
le = latest.get((t, p), -1)
if off < 0 or le < 0:
lag_str, status = "?", "no-data"
else:
lag = max(0, le - off)
total_lag += lag
lag_str = str(lag)
status = "OK" if lag == 0 else (
"⚠ growing" if lag > 1000 else "small")
print(f" {t:<26}{p:<6}{off:<14}{le:<14}{lag_str:<10}{status}")
print(f" {'TOTAL LAG':<46} {total_lag}")
def main():
args = parse_args()
admin = AdminClient({"bootstrap.servers": args.bootstrap})
while True:
for group in args.group:
try:
state, members, committed = fetch_state(admin, group)
tps = [TopicPartition(t, p) for (t, p) in committed.keys()]
latest = fetch_latest_offsets(admin, tps)
render(group, state, members, committed, latest, args)
except Exception as e:
print(f"[ERR] group={group} {e}", file=sys.stderr)
if args.once:
return
time.sleep(args.interval)
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
sys.exit(0)python
#!/usr/bin/env python3
"""
rebalance_listener.py — 完整演示 Rebalance 监听器的「资源清理」范式
核心范式(生产环境推荐):
1) on_assign 里:分配到的每个分区,初始化它的本地状态
(如:去重布隆过滤器、聚合缓冲区、外部连接句柄);
2) on_revoke 里:要失去的分区,先 flush 业务缓冲、同步 commit offset、
清理本地状态——保证下一个接手的 Consumer 不会重做或丢失工作;
3) on_lost 里(仅 cooperative):分区已经被强制收走,**不要 commit**;
只清理本地状态。
实验场景:
本脚本启动一个 Consumer,订阅 Topic learn.11.listener。
模拟「每个分区维护一个本地累加计数器 + 一个 dedup set」。
在 Rebalance 时打印:保存了多少状态、commit 了哪个 offset。
建议同时启动 2~3 个本程序的实例,观察分配/撤销过程。
用法:
python rebalance_listener.py --consumer-id c1
"""
from __future__ import annotations
import argparse
import signal
import sys
import time
from collections import defaultdict
from datetime import datetime
from confluent_kafka import Consumer, KafkaError, TopicPartition
def parse_args():
p = argparse.ArgumentParser()
p.add_argument("--bootstrap", default="localhost:9092")
p.add_argument("--topic", default="learn.11.listener")
p.add_argument("--group", default="listener-demo-group")
p.add_argument("--consumer-id", default="c1")
return p.parse_args()
class PartitionState:
"""每个分区的本地业务状态"""
__slots__ = ("counter", "dedup", "last_offset", "buffered_writes")
def __init__(self):
self.counter = 0 # 累计消费条数
self.dedup = set() # 业务侧去重 set(模拟)
self.last_offset = -1 # 最近处理过的 offset
self.buffered_writes = 0 # 还没 flush 的写入条数
class RebalanceHandler:
"""统一封装 on_assign / on_revoke / on_lost 与本地状态管理"""
def __init__(self, consumer_id: str):
self.cid = consumer_id
self.states: dict[tuple[str, int], PartitionState] = {}
def _log(self, msg, color="\033[36m"):
print(f"{color}[{datetime.now().strftime('%H:%M:%S')}][{self.cid}] {msg}\033[0m",
flush=True)
def on_assign(self, consumer: Consumer, partitions):
self._log(f"on_assign +{len(partitions)} parts: "
f"{[(tp.topic, tp.partition) for tp in partitions]}", "\033[32m")
for tp in partitions:
key = (tp.topic, tp.partition)
if key not in self.states:
self.states[key] = PartitionState()
self._log(f" ↳ init local state for {key}", "\033[32m")
consumer.incremental_assign(partitions)
def on_revoke(self, consumer: Consumer, partitions):
self._log(f"on_revoke -{len(partitions)} parts: "
f"{[(tp.topic, tp.partition) for tp in partitions]}", "\033[33m")
# 1. flush 业务侧未持久化的写入
for tp in partitions:
key = (tp.topic, tp.partition)
st = self.states.get(key)
if st and st.buffered_writes > 0:
self._log(f" ↳ flush {st.buffered_writes} buffered writes "
f"for {key}", "\033[33m")
# 这里假装 flush 到 DB
st.buffered_writes = 0
# 2. 同步 commit 当前 offset,避免重复消费
try:
offsets_to_commit = []
for tp in partitions:
key = (tp.topic, tp.partition)
st = self.states.get(key)
if st and st.last_offset >= 0:
# commit 的 offset 是「下一条要消费的」,所以 +1
offsets_to_commit.append(
TopicPartition(tp.topic, tp.partition, st.last_offset + 1))
if offsets_to_commit:
consumer.commit(offsets=offsets_to_commit, asynchronous=False)
self._log(f" ↳ sync committed offsets: "
f"{[(o.topic, o.partition, o.offset) for o in offsets_to_commit]}",
"\033[33m")
except Exception as e:
self._log(f" ↳ commit failed: {e}", "\033[31m")
# 3. 释放本地状态
for tp in partitions:
key = (tp.topic, tp.partition)
st = self.states.pop(key, None)
if st:
self._log(f" ↳ released local state for {key} "
f"(processed={st.counter}, dedup_size={len(st.dedup)})",
"\033[33m")
consumer.incremental_unassign(partitions)
def on_lost(self, consumer: Consumer, partitions):
# cooperative 模式独有:分区已经被强制收回,不能 commit
self._log(f"on_lost -{len(partitions)} parts (forced, NO commit)",
"\033[31m")
for tp in partitions:
key = (tp.topic, tp.partition)
self.states.pop(key, None)
consumer.incremental_unassign(partitions)
def process(self, msg):
key = (msg.topic(), msg.partition())
st = self.states.get(key)
if not st:
# Rebalance 中刚被 revoke,丢弃
return
st.counter += 1
st.last_offset = msg.offset()
st.buffered_writes += 1
# 模拟去重 key(这里用 offset 作为 dedup id)
st.dedup.add(f"{msg.topic()}-{msg.partition()}-{msg.offset()}")
if st.counter % 100 == 0:
# 模拟周期性 flush
self._log(f" periodic flush: {key} -> {st.buffered_writes} writes",
"\033[34m")
st.buffered_writes = 0
def main():
args = parse_args()
handler = RebalanceHandler(args.consumer_id)
consumer = Consumer({
"bootstrap.servers": args.bootstrap,
"group.id": args.group,
"client.id": args.consumer_id,
"partition.assignment.strategy": "cooperative-sticky",
"enable.auto.commit": False,
"auto.offset.reset": "latest",
"session.timeout.ms": 45000,
"heartbeat.interval.ms": 3000,
"max.poll.interval.ms": 300000,
})
consumer.subscribe([args.topic],
on_assign=handler.on_assign,
on_revoke=handler.on_revoke,
on_lost=handler.on_lost)
handler._log(f"started, subscribing to {args.topic}", "\033[35m")
stop = {"flag": False}
def _sig(*_):
handler._log("got signal, shutting down…", "\033[35m")
stop["flag"] = True
signal.signal(signal.SIGINT, _sig)
signal.signal(signal.SIGTERM, _sig)
last_print = time.time()
try:
while not stop["flag"]:
msg = consumer.poll(1.0)
if msg is None:
if time.time() - last_print > 5:
last_print = time.time()
held = consumer.assignment()
handler._log(
f"holding {len(held)} partitions, "
f"states={list(handler.states.keys())}",
"\033[34m")
continue
if msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
continue
handler._log(f"err: {msg.error()}", "\033[31m")
continue
handler.process(msg)
finally:
# 优雅关闭:close() 会触发 on_revoke 路径
handler._log("draining and closing…", "\033[35m")
consumer.close()
handler._log("closed cleanly.", "\033[35m")
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
sys.exit(0)python
#!/usr/bin/env python3
"""
static_membership.py — 演示 group.instance.id(Static Membership, KIP-345)
实验流程(建议同时开 4 个终端):
终端 1(启用 static membership 的 Consumer A):
python static_membership.py --consumer-id static-a --instance-id pod-a
终端 2(启用 static membership 的 Consumer B):
python static_membership.py --consumer-id static-b --instance-id pod-b
终端 3(启用 static membership 的 Consumer C):
python static_membership.py --consumer-id static-c --instance-id pod-c
然后在终端 4 用 kafka-consumer-groups.sh 看消费组状态:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group static-mem-demo --state
实验:CTRL+C 关掉终端 1(pod-a),10 秒内重启它:
* 在 60s session.timeout.ms 内,Coordinator 不会触发 Rebalance;
* pod-a 重新启动后,原来的分区直接还回来,整组无 Rebalance;
* 终端 2 / 3 完全不停。
对比试验:去掉 --instance-id 参数(变成 dynamic membership)再重做:
* 关掉某个 Consumer 立刻触发 Rebalance;
* 启动它又触发一次。整组停顿明显。
注意:
* group.instance.id 必须**全局唯一**;如果你启动两个相同的 instance.id,
后启动的会顶替先启动的(先启动的会拿到 fenced 错误);
* static membership 要求 broker 端 inter.broker.protocol.version >= 2.3。
"""
from __future__ import annotations
import argparse
import signal
import sys
import time
from datetime import datetime
from confluent_kafka import Consumer
def parse_args():
p = argparse.ArgumentParser()
p.add_argument("--bootstrap", default="localhost:9092")
p.add_argument("--topic", default="learn.11.static")
p.add_argument("--group", default="static-mem-demo")
p.add_argument("--consumer-id", default="c1", help="client.id")
p.add_argument("--instance-id", default=None,
help="group.instance.id;不传则退化为 dynamic membership")
p.add_argument("--session-ms", type=int, default=60000,
help="session.timeout.ms(static 模式建议调到 30~60s)")
return p.parse_args()
def fmt_tps(tps):
return ", ".join(f"{tp.topic}-{tp.partition}" for tp in tps) or "<empty>"
def main():
args = parse_args()
cid = args.consumer_id
iid = args.instance_id
color = "\033[32m" if iid else "\033[33m"
tag = f"static:{iid}" if iid else f"dynamic:{cid}"
def log(msg, c=color):
print(f"{c}[{datetime.now().strftime('%H:%M:%S')}][{tag}] {msg}\033[0m",
flush=True)
cfg = {
"bootstrap.servers": args.bootstrap,
"group.id": args.group,
"client.id": cid,
"partition.assignment.strategy": "cooperative-sticky",
"enable.auto.commit": False,
"auto.offset.reset": "latest",
"session.timeout.ms": args.session_ms,
"heartbeat.interval.ms": max(3000, args.session_ms // 5),
"max.poll.interval.ms": 300000,
}
if iid:
# ★关键★ 启用 static membership
cfg["group.instance.id"] = iid
def on_assign(consumer, partitions):
log(f"on_assign → +{fmt_tps(partitions)}", "\033[36m")
consumer.incremental_assign(partitions)
def on_revoke(consumer, partitions):
log(f"on_revoke → -{fmt_tps(partitions)}", "\033[35m")
try:
consumer.commit(asynchronous=False)
except Exception as e:
log(f"commit failed: {e}", "\033[31m")
consumer.incremental_unassign(partitions)
consumer = Consumer(cfg)
consumer.subscribe([args.topic], on_assign=on_assign, on_revoke=on_revoke)
log(f"started. cfg.session_ms={args.session_ms} "
f"static={'YES' if iid else 'NO'}")
log("提示:CTRL+C 关掉本进程,10s 内 rerun,观察是否触发 Rebalance")
stop = {"flag": False}
def _sig(*_):
stop["flag"] = True
signal.signal(signal.SIGINT, _sig)
signal.signal(signal.SIGTERM, _sig)
msgs = 0
last_print = time.time()
try:
while not stop["flag"]:
msg = consumer.poll(1.0)
if msg is None:
if time.time() - last_print > 5:
last_print = time.time()
held = consumer.assignment()
log(f"alive | holding: {fmt_tps(held)} | msgs={msgs}", "\033[34m")
continue
if msg.error():
log(f"err: {msg.error()}", "\033[31m")
continue
msgs += 1
if msgs <= 3 or msgs % 50 == 0:
log(f"poll {msg.topic()}-{msg.partition()}@{msg.offset()}", "\033[36m")
consumer.commit(message=msg, asynchronous=True)
finally:
log("closing…", "\033[35m")
# 优雅 close 会发 LeaveGroup
# 但 static membership 的 close() 行为可以通过 leave-group-on-close 控制
# 此处保持默认(dynamic 会发 LeaveGroup,static 不发)
consumer.close()
log(f"closed. msgs={msgs}", "\033[35m")
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
sys.exit(0)cooperative_consumer.py ↗ · lag_monitor.py ↗ · rebalance_listener.py ↗ · static_membership.py ↗