现象:每次发版,消费停 47 秒
9 月底做容量复盘时,我发现一个奇怪的规律:每次发布,Kafka 消费 lag 都会先涨后落,中间有大约 45 到 60 秒消费完全停住。看消费者日志,那段时间全是这个:
2021-09-28 14:22:31.407 WARN [Consumer clientId=consumer-order-3, groupId=order-group]
o.a.k.c.c.internals.ConsumerCoordinator : [Consumer clientId=consumer-order-3,
groupId=order-group] Attempt to heartbeat failed since group is rebalancing
2021-09-28 14:22:33.512 INFO o.a.k.c.c.internals.ConsumerCoordinator :
[Consumer clientId=consumer-order-3, groupId=order-group] Revoking previously
assigned partitions [ORDER_PAY-3, ORDER_PAY-7, ORDER_PAY-11, ...]
2021-09-28 14:23:18.774 INFO o.a.k.c.c.internals.ConsumerCoordinator :
[Consumer clientId=consumer-order-3, groupId=order-group] Setting newly
assigned partitions [ORDER_PAY-3, ORDER_PAY-7, ORDER_PAY-11, ...]
从 14:22:33 撤销分区到 14:23:18 重新分配,45 秒消费空窗。我们滚动发布 12 个实例,每次实例重启都触发一次全组 rebalance,等于发了 12 次,全线停了将近 9 分钟。
先搞清楚 rebalance 什么时候会触发
消费者组的重平衡由 GroupCoordinator 驱动,触发条件就这几种:
- 组成员数量变化:新消费者加入、已有消费者正常退出、或者崩溃(心跳超时)。这是最常见的一种,我们遇到的是这个。
- 订阅 topic 的分区数增加:
kafka-topics.sh --alter --partitions之后,所有订阅者要重新分配。 - 订阅的 topic 集合变化。用
consumer.subscribe(Pattern)正则订阅时,新建了匹配的 topic 就会触发。我们有一次就是运维建了个ORDER_PAY_TEST,正好被正则匹配到,把生产组搅了一遍。 - 消费者超过
max.poll.interval.ms没发起下一次 poll,被判定为"僵死"踢出组。这个很隐蔽,处理批消息耗时太长就会中招。 - 消费者超过
session.timeout.ms没发心跳:网络抖动、或者一次长 GC STW。
第 4、5 条值得单独说。我们之前把 max.poll.records 设成 2000,批量处理平均要 12 秒,某次下游抖动处理了 6 分钟,直接超过 max.poll.interval.ms 的 5 分钟默认值,消费者被踢出组,然后又重新加入,来回震荡。
eager 协议:为什么要停全世界
Kafka 2.x 默认用的是 eager(急切)重平衡协议。它的流程是:
- Coordinator 检测到成员变化,通知所有成员"要重平衡了"。
- 所有成员必须撤销自己持有的全部分区,提交 offset,然后重新加入组。
- Coordinator 等所有成员都重新加入(或者等
rebalance.timeout.ms,默认 5 分钟),期间消费完全停止。 - 选出 leader 消费者,由它按分配策略算出方案,广播给所有人。
- 所有人重新开始消费。
关键在第 2 步:哪怕只有一个分区需要移动,全组所有分区都要先撤销。这是设计上的简化,代价就是全量停顿。分区越多、消费者越多,停顿越长。
我们那个组订阅了 4 个 topic 共 96 个分区,12 个消费者。每次 rebalance 要:撤销 96 个分区的 offset 提交 + 全组 12 个成员重新加入的协调开销 + 重新分配广播。45 秒就是这么来的。
用命令行能直接看到组的状态:
$ kafka-consumer-groups.sh --bootstrap-server 10.0.1.20:9092 \
--group order-group --describe --state
GROUP COORDINATOR ASSIGNMENT-STRATEGY STATE #MEMBERS
order-group 10.0.1.20:9092 (id: 2147483644) range PreparingRebalance 8
STATE 在 PreparingRebalance / CompletingRebalance 之间切换时就是停顿时。Stable 才是正常。
改法一:静态成员(KIP-345)
发版导致的 rebalance 完全是可以避免的。Kafka 2.3 引入了静态成员(KIP-345),给每个消费者配一个固定的 group.instance.id:
spring:
kafka:
consumer:
group-id: order-group
# 关键:每个实例一个固定的、重启后不变的 ID
properties:
group.instance.id: ${POD_NAME}
POD_NAME 从 K8s 的 Downward API 注入:
env:
- name: POD_NAME
valueFrom:
fieldRef:
fieldPath: metadata.name
原理:Coordinator 认的是 group.instance.id 而不是每次连接生成的 member.id。实例重启时,消费者带着同样的 instance.id 重新加入,Coordinator 认为"还是那个老成员",只要它在 session.timeout.ms 内回来,就不触发 rebalance,直接把原来的分区还给它。
这就要求我们把 session.timeout.ms 调到能覆盖一次重启的时间。我们应用从收到 SIGTERM 到新 Pod 就绪大约 35 秒(含 Spring 启动 25 秒),所以设成 60 秒:
spring:
kafka:
consumer:
properties:
group.instance.id: ${POD_NAME}
session.timeout.ms: 60000 # 默认 45000,改大
heartbeat.interval.ms: 3000 # 保持 session.timeout 的 1/20 左右
max.poll.interval.ms: 300000
max.poll.records: 500 # 从 2000 降下来
注意:静态成员的代价是"实例真挂了要等 session.timeout.ms 才恢复"。以前崩溃 10 秒就被踢出组、分区立刻转给别人;现在要等 60 秒。这 60 秒里那些分区是没人消费的。所以 session.timeout.ms 不能无脑调大,要按实际的重启耗时来定(我们的 35 秒 + 25% 余量)。
另外滚动发布时要保证"先起新的,再停老的"会造成成员数短暂超标。K8s 默认的 RollingUpdate 是 maxSurge=25%, maxUnavailable=0,会先起新 Pod。有静态成员的情况下,新 Pod(新 instance.id 如果 POD_NAME 变了)会被当成新成员加入,触发一次 rebalance。我们的处理是把 Deployment 改成 maxSurge: 0,先停旧的再起新的,配合静态成员就没有多余 rebalance。
改法二:增量协作式重平衡(KIP-429)
静态成员解决了"计划内重启",但分区数变化、成员真的增减这些还是得 rebalance。Kafka 2.4 引入的协作式重平衡(KIP-429)能大幅缩短停顿。
思路和 eager 的关键区别:eager 是一次性撤销全部分区,cooperative 分两轮,只撤销真正需要移动的分区。
- 第一轮:Coordinator 通知大家"要重平衡"。所有人继续消费,只是把需要转让的分区标记为待撤销。
- 第二轮:Coordinator 把第一轮收集到的结果汇总,把待撤销的分区分配给新成员。没被移动的分区全程没有中断。
启用方式是换分配策略:
spring:
kafka:
consumer:
properties:
partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignor
CooperativeStickyAssignor 是"协作式 + 粘性"的组合:协作式保证分区尽量不中断,粘性保证重平衡后尽量维持原分配(减少分区在实例间来回迁移,也就减少了每个实例重新预热本地缓存的成本)。
升级到 cooperative 协议时,滚动发布期间会同时存在新旧两种协议的成员。Kafka 的处理是:只要还有一个老成员(eager 协议),整个组就降级走 eager。所以必须等所有实例都发完新版,cooperative 才真正生效。这个过渡期要注意观察日志里的
protocol相关提示。
优雅退出也是减少 rebalance 的一环
还有一个容易被忽略的点:进程收到 SIGTERM 之后,如果没有主动关闭消费者,Coordinator 要等 session.timeout.ms 才发现"这个人走了",期间会触发一次 rebalance。正确的做法是收到退出信号时主动 close() 消费者,它会发 LeaveGroup 请求,Coordinator 立刻把分区转给别人。
Spring Kafka 里默认是有这个逻辑的(KafkaListenerEndpointRegistry 会在容器销毁时 stop 监听容器),但要保证关闭等待时间够长。我们一开始只给了 10 秒:
spring:
kafka:
listener:
# 容器关闭时等待消费中的消息处理完的时间,默认 10 秒
shutdown-timeout: 30000
批量消费一批 500 条要 90 秒,10 秒根本不够,导致每次发版都会有一批消息没处理完就被中断,重启后又从头消费(offset 没提交)。调到 30 秒并配合 max.poll.records=500 才对齐。
K8s 侧也要配套,给足优雅退出时间:
spec:
terminationGracePeriodSeconds: 60 # 默认 30 秒,我们改成 60
containers:
- lifecycle:
preStop:
exec:
command: ["sh", "-c", "sleep 5"] # 给 endpoint 摘除留时间
还可以在关闭前先暂停消费,让在途消息处理干净:
@Autowired
private KafkaListenerEndpointRegistry registry;
public void gracefulShutdown() {
registry.getAllListenerContainers()
.forEach(MessageListenerContainer::pause); // 停止拉取新消息
// 等待在途消息处理完
Thread.sleep(20_000);
registry.destroy(); // 真正关闭
}
我们做了什么
两步都做了,还顺带调整了几个参数:
spring:
kafka:
listener:
ack-mode: BATCH
concurrency: 8 # 每个实例的消费线程数,和分区数匹配
consumer:
group-id: order-group
max-poll-records: 500 # 从 2000 降下来
properties:
group.instance.id: ${POD_NAME}
session.timeout.ms: 60000
heartbeat.interval.ms: 3000
max.poll.interval.ms: 300000
partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignor
# 2.6 起的心跳优化,让 rebalance 检测更快
reconnect.backoff.max.ms: 1000
消费逻辑也改了:max.poll.records 从 2000 降到 500,单批耗时从最长 6 分钟降到 90 秒以内,远离 max.poll.interval.ms 的 5 分钟红线。并且把批处理里的外部 HTTP 调用加了降级开关(和 MQ 那篇一样的思路)。
效果
| 指标 | 优化前 | 优化后 |
|---|---|---|
| 单次发版的消费停顿 | 45 ~ 62 秒 | 0.8 ~ 1.4 秒 |
| 日均 rebalance 次数 | 341 次 | 3 次 |
| rebalance 期间最大 lag | 186 万 | 1.2 万 |
| 分区分配策略 | range | CooperativeStickyAssignor |
| 因 max.poll.interval 被踢出 | 每周 4 ~ 7 次 | 0 |
剩下那 3 次 rebalance 是分区扩容和真的实例故障,属于正常情况。停顿从 45 秒降到 1 秒左右,因为需要移动的分区极少(静态成员把绝大部分分区固定在原实例上了)。
监控得跟上
光改不监控等于没改。Kafka 消费者暴露了几个 JMX 指标,我们接到了 Prometheus:
# 重平衡频率,正常情况下应该接近 0
rate(kafka_consumer_rebalance_total{group="order-group"}[10m])
# 单次重平衡耗时
kafka_consumer_rebalance_time_avg_seconds{group="order-group"}
# 上次重平衡的耗时(秒),这个更直观
kafka_consumer_last_rebalance_seconds{group="order-group"}
# 组内成员数,突然变化说明有实例掉线
kafka_consumer_group_member_count{group="order-group"}
告警规则:
- alert: KafkaRebalanceTooFrequent
expr: increase(kafka_consumer_rebalance_total{group="order-group"}[30m]) > 5
for: 5m
labels:
severity: warning
annotations:
summary: "{{ $labels.group }} 30 分钟内重平衡 {{ $value }} 次"
- alert: KafkaRebalanceSlow
expr: kafka_consumer_last_rebalance_seconds{group="order-group"} > 10
for: 2m
labels:
severity: warning
annotations:
summary: "单次重平衡耗时 {{ $value }} 秒,消费中断过久"
我们的 spring-kafka 版本是 2.7.x(对应 kafka-clients 2.7.1)。顺带提一句版本:Kafka 3.0 在 2021 年 9 月发布,客户端和 broker 都要求 Java 11+,2.8 是最后一个支持 Java 8 的版本。我们当时还在 JDK 8 和 11 混合期,所以锁在 2.8.x,没急着升 3.0。
小结
- rebalance 的五种触发条件:成员增减、分区数变化、订阅 topic 变化、
max.poll.interval.ms超时、session.timeout.ms超时。后两种最容易被忽略,处理批量消息耗时长就会中招。 - eager 协议的问题是"撤销全部分区",哪怕只动一个也要全停,这是 45 秒停顿的根源。
- 静态成员(
group.instance.id)解决计划内重启,K8s 环境用 Downward API 注入POD_NAME即可。代价是真故障时要等session.timeout.ms才恢复,这个值要按实际重启耗时定,别无脑调大。 - 增量协作式重平衡(
CooperativeStickyAssignor)解决计划外的分区迁移,只中断需要移动的分区。注意滚动发布期间新旧协议共存会降级为 eager,要等全量发完才生效。 - 降低
max.poll.records、给批量处理加降级开关,能避免消费者被误判为僵死。 - 监控
kafka_consumer_rebalance_total和kafka_consumer_last_rebalance_seconds,这两个指标比看日志直观得多。 - 版本上注意:Kafka 2.8 是最后一个支持 Java 8 的版本,3.0 起要求 Java 11+。