Skip to content

第 11 章 消费者组与 Rebalance:Kafka 客户端「最大的痛」

目标读者:写过 Consumer 但说不清「Rebalance 到底什么时候触发」「为什么我的消费在动态扩容时全停了 30 秒」「group.instance.id 是什么」「max.poll.interval.mssession.timeout.ms 哪个先生效」的同学。

学完你会:能画出 Consumer 加入 Group 的完整 4 步握手,能说清 Range/RoundRobin/Sticky/CooperativeSticky 的差别,知道 Eager Rebalance 为什么会「全停」、Cooperative 为什么能「增量」,能用 group.instance.id 治住临时网络抖动引发的 Rebalance 风暴,遇到 Lag 持续上涨能按清单逐项排查。


0. 导读:消费者组是 Kafka 「最朴素却最容易踩坑」的设计

很多同学觉得 Consumer 写起来比 Producer 简单:subscribepoll、处理消息、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

也就是说:

  1. Kafka 内部有个 Topic __consumer_offsets,默认 50 个分区;
  2. 把消费组名 hash,映射到某个分区;
  3. 那个分区的 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 四个请求各自做什么

请求谁发谁回干什么
FindCoordinatorConsumer任意 Broker通过 hash(group.id) 找到 Coordinator 的 broker
JoinGroupConsumerGroup Coordinator「我要加入这个组」,注册自己支持的分配策略列表
SyncGroupConsumer (Group Leader)Group CoordinatorGroup Leader 把计算好的分配方案告知 Coordinator 转发
HeartbeatConsumerGroup Coordinator「我还活着」,Coordinator 也借此通知有 Rebalance 发生
LeaveGroupConsumerGroup Coordinator优雅下线,触发 Rebalance(不发也行,但要等心跳超时)
OffsetCommitConsumerGroup Coordinator提交 offset 到 __consumer_offsets Topic
OffsetFetchConsumerGroup Coordinator启动时拉自己负责的分区上一次的 offset

2.3 Group Leader vs Coordinator(容易混)

  • Group CoordinatorBroker 端的角色——某台 Broker 兼任,负责接受请求、维护状态。
  • Group LeaderConsumer 端的角色——同一个 Group 中第一个 join 进来的 Consumer 被指定为 Leader,它在客户端本地计算「谁分到哪些分区」,然后把方案通过 SyncGroup 上报给 Coordinator。

为什么把分配算法放在客户端? 这是 Kafka 一个聪明的设计:分配算法可以独立演进、用户可以自定义、Broker 端不需要理解协议细节。Broker 只是个透传方案的「邮局」。


3. 触发 Rebalance 的所有情形

记住一个总规则:消费组成员构成 / 订阅关系 / 分区数发生变化时,都会触发 Rebalance。具体:

  1. 新 Consumer 加入 Group(启动 / 扩容);
  2. Consumer 离开 Group(主动 LeaveGroup / 进程崩溃 / 心跳超时 / 处理超时);
  3. 订阅的 Topic 变化subscribe(["a", "b"]) 改成 subscribe(["a", "b", "c"]));
  4. 订阅的 Topic 分区数变化(管理员扩了分区);
  5. 正则订阅 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 里了」,导致:

  1. 整组 Rebalance 30 秒(其它 Consumer 也跟着停);
  2. 你处理完的那条消息可能要被另一个 Consumer 重消费(如果你还没 commit);
  3. 你的 Consumer 重新 join 又触发一次 Rebalance;
  4. 业务逻辑慢的根因没解决,下一批消息又会超时——死循环 Rebalance 风暴

修复思路:调大 max.poll.interval.ms 是治标,真正的根治是把每条消息处理时间砍下来——拆成异步、加并发、降批量。


4. 四大分区分配策略对比

Kafka 内置 4 种分配策略(partition.assignment.strategy):

  1. RangeAssignor(默认,单 Topic 范围分)
  2. RoundRobinAssignor(轮询)
  3. StickyAssignor(粘性,2.4 之前)
  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)做了两件事:

  1. 均匀(与 RoundRobin 一样均匀,不一定一致);
  2. 粘性: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 黄金法则

  1. on_revoke 一定要 commit(同步 commit,不能 fire-and-forget),否则同一条消息会被新接手的 Consumer 重消费;
  2. on_revoke 是阻塞的:你在里面做的任何事都会延长 Rebalance 时间;不要在里面调外部接口、做大量 IO;
  3. on_lost 不要 commit:你已经没有这些分区的所有权,commit 会被拒;
  4. 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 --execute

8. 配置速查表

8.1 客户端必配

配置默认推荐说明
group.id-业务名消费组标识
partition.assignment.strategyrangecooperative-sticky强烈推荐
enable.auto.committruefalse手动 commit 才能保证不丢
auto.offset.resetlatestearliest(首次/重置)没有 commit 时从哪开始
session.timeout.ms4500030000~60000配 static membership 时调大
heartbeat.interval.ms3000session.timeout / 3Heartbeat 频率
max.poll.interval.ms300000看业务最长处理时间业务慢就调大或拆消费
max.poll.records500视消息大小一次 poll 最多拿多少条
fetch.min.bytes / max.wait.ms1 / 500高吞吐场景调大Broker 端聚合等待
group.instance.id-用 Pod 名静态成员,治抖

8.2 服务端相关

配置默认说明
offsets.topic.replication.factor3__consumer_offsets 副本数(务必 ≥3)
offsets.topic.num.partitions50分区数(决定能承载多少消费组并发)
group.initial.rebalance.delay.ms3000初始 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 条:

  1. 首选 cooperative-sticky 策略:从 Kafka 2.4 起就有,能省 80% 的「全停时间」;
  2. group.instance.id + 调大 session.timeout.ms:治住临时抖动引发的「假性 Rebalance」;
  3. 业务处理时间不要超过 max.poll.interval.ms:宁可拆细每条消息的处理,也别让 poll 超时;
  4. 手动 commit + on_revoke 同步 commit:防止重复消费;
  5. 消费 Lag 接监控告警:盯住 LAG 与 STATE 两个指标,频繁 Rebalance 比 Lag 上涨更可怕。

11. 配套代码与演示

  • 11_consumer_group/code/cooperative_consumer.pypartition.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 节点宕机:
    1. __consumer_offsets 那个分区会经历 ISR 切换,新 Leader 上线;
    2. Consumer 发请求到旧 Coordinator 收到 NOT_COORDINATOR 错误;
    3. Consumer 重新调用 FindCoordinator 定位到新 Coordinator;
    4. Consumer 重连后发 JoinGroup,触发一次 Rebalance;
  • 所以 __consumer_offsets 的副本数(offsets.topic.replication.factor)一定要 ≥ 3。

Q2:详细描述一个 Consumer 加入 Group 的全过程

考察点:协议层的细节理解。

标准答案:四步:

  1. FindCoordinator:Consumer 拿 group.id 询问任意 Broker,返回 Coordinator 所在 Broker;
  2. JoinGroup:Consumer 向 Coordinator 发 JoinGroup(带支持的策略列表),Coordinator 等够 group.initial.rebalance.delay.ms 后把全员状态返回;返回时指定第一个 join 的 Consumer 为 Group Leader
  3. SyncGroup:Group Leader 在客户端本地按选定策略计算分配方案,把方案随 SyncGroup 上报;其它 Consumer 也发 SyncGroup(不带方案),Coordinator 收齐后把方案分发给每个 Consumer;
  4. 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.msmax.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 持续上涨怎么排查?

考察点:实战能力。

标准答案(按概率从高到低):

  1. 生产>消费:看 Producer QPS 和 Consumer QPS 对比,扩 Consumer 实例、调大 max.poll.records;
  2. 业务处理太慢导致 max.poll.interval.ms 超时:看 server.log 是否高频 Rebalance;改进消息处理逻辑;
  3. 单 Consumer 实例打满:看 CPU/GC/IO;
  4. Key 热点导致部分分区积压:用 --describe 看每分区 LAG;
  5. CONSUMER-ID 是 -:根本没人在消费;
  6. 消费组在 Rebalance:看 --state,长期非 Stable 要查原因;
  7. Coordinator 节点故障:检查 __consumer_offsets 对应分区的 Leader Broker;
  8. Topic 加了分区但 Consumer 没感知:等 metadata.max.age.ms 或重启 Consumer。

排查工具kafka-consumer-groups.sh --describe、JMX records-lag-max、Kafka UI。

Q7:自动 commit 和手动 commit 各有什么风险?

考察点:消息语义。

标准答案

  • 自动 commitenable.auto.commit=true,默认每 5s 提交一次):
    • 风险 1(重复消费):处理了消息但还没到自动提交的时刻就崩溃,重启会从上一次自动 commit 的位置开始 → 重复消费;
    • 风险 2(消息丢失):自动 commit 是「按上次 poll 拿到的最大 offset」提交,poll 之后还在处理,commit 已经过去了——如果业务处理失败但 commit 已发,那条消息就丢了;
  • 手动 commit
    • 同步 commit (commitSync):阻塞,可靠但慢;
    • 异步 commit (commitAsync):不阻塞,但出错没回调就忘了;
    • 推荐组合:日常 commitAsync + 关闭/Rebalance 时 commitSync + on_revoke 里同步 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 会周期性地刷新 metadatametadata.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 ↗