Administrator
发布于 2021-10-15 / 3285 阅读
73

Kafka 消费者组重平衡问题与优化

现象:每次发版,消费停 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 驱动,触发条件就这几种:

  1. 组成员数量变化:新消费者加入、已有消费者正常退出、或者崩溃(心跳超时)。这是最常见的一种,我们遇到的是这个。
  2. 订阅 topic 的分区数增加kafka-topics.sh --alter --partitions 之后,所有订阅者要重新分配。
  3. 订阅的 topic 集合变化。用 consumer.subscribe(Pattern) 正则订阅时,新建了匹配的 topic 就会触发。我们有一次就是运维建了个 ORDER_PAY_TEST,正好被正则匹配到,把生产组搅了一遍。
  4. 消费者超过 max.poll.interval.ms 没发起下一次 poll,被判定为"僵死"踢出组。这个很隐蔽,处理批消息耗时太长就会中招。
  5. 消费者超过 session.timeout.ms 没发心跳:网络抖动、或者一次长 GC STW。

第 4、5 条值得单独说。我们之前把 max.poll.records 设成 2000,批量处理平均要 12 秒,某次下游抖动处理了 6 分钟,直接超过 max.poll.interval.ms 的 5 分钟默认值,消费者被踢出组,然后又重新加入,来回震荡。

eager 协议:为什么要停全世界

Kafka 2.x 默认用的是 eager(急切)重平衡协议。它的流程是:

  1. Coordinator 检测到成员变化,通知所有成员"要重平衡了"。
  2. 所有成员必须撤销自己持有的全部分区,提交 offset,然后重新加入组。
  3. Coordinator 等所有成员都重新加入(或者等 rebalance.timeout.ms,默认 5 分钟),期间消费完全停止。
  4. 选出 leader 消费者,由它按分配策略算出方案,广播给所有人。
  5. 所有人重新开始消费。

关键在第 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

STATEPreparingRebalance / 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 分两轮,只撤销真正需要移动的分区

  1. 第一轮:Coordinator 通知大家"要重平衡"。所有人继续消费,只是把需要转让的分区标记为待撤销。
  2. 第二轮: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 期间最大 lag186 万1.2 万
分区分配策略rangeCooperativeStickyAssignor
因 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_totalkafka_consumer_last_rebalance_seconds,这两个指标比看日志直观得多。
  • 版本上注意:Kafka 2.8 是最后一个支持 Java 8 的版本,3.0 起要求 Java 11+。

参考