Administrator
发布于 2020-06-02 / 27614 阅读
341

Kafka 消费积压的排查与扩容处理

早上八点,消费积压告警:lag 320 万

6 月 1 号早上八点十几分,钉钉机器人开始刷屏:ORDER_EVENT_TOPIC 的消费 lag 突破 300 万,还在涨。

这个 topic 是订单事件流,下游有 5 个消费者组:风控、数据同步、搜索索引、积分、客服。lag 涨到 320 万意味着这些下游全部延迟,最严重的是搜索索引——商品搜索结果要 40 分钟才能更新。

从排查到恢复用了 2 小时 20 分。这篇把过程、原理和之后做的优化写下来。

先看 lag 长什么样

Kafka 2.4.1,用自带的脚本看消费者组状态:

$ bin/kafka-consumer-groups.sh --bootstrap-server 10.0.0.41:9092 \
    --describe --group search-index-group

Consumer group 'search-index-group' has no active members.

TOPIC              PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG     CONSUMER-ID  HOST
ORDER_EVENT_TOPIC  0          8412034         8672911         260877  -            -
ORDER_EVENT_TOPIC  1          8402117         8670338         268221  -            -
ORDER_EVENT_TOPIC  2          8398002         8668912         270910  -            -
...
ORDER_EVENT_TOPIC  11         8411990         8671405         259415  -            -

第一行的 has no active members 是关键信息:消费者组里一个活着的成员都没有。不是消费慢,是完全停了。

LAG 是 26 万到 27 万,12 个分区,加起来 320 万。每个分区的 lag 差不多,说明是均匀停摆而不是某个分区卡住。

为什么消费者全没了

去看消费者服务的日志,找到了这个:

2020-06-01 07:41:22.318  WARN  [search-index-0] o.a.k.c.c.internals.ConsumerCoordinator
  : [Consumer clientId=consumer-1, groupId=search-index-group] This member will leave the group
    because consumer poll timeout has expired. This means the time between subsequent calls to
    poll() was longer than the configured max.poll.interval.ms, which typically implies that the
    poll loop is spending too much time processing messages. You can address this either by
    increasing max.poll.interval.ms or by reducing the max size of batches returned in poll()
    with max.poll.records.

2020-06-01 07:41:22.401  INFO  [search-index-0] o.a.k.c.c.internals.ConsumerCoordinator
  : [Consumer clientId=consumer-1, groupId=search-index-group] (Re-)joining group

原因很清楚:max.poll.interval.ms 超时了。这个参数的含义是"两次 poll() 调用之间的最大间隔",超过就认为这个消费者挂了,踢出组、触发 rebalance。

我们的消费者代码长这样:

@KafkaListener(topics = "ORDER_EVENT_TOPIC", groupId = "search-index-group",
               concurrency = "3")
public void consume(List<ConsumerRecord<String, String>> records) {
    for (ConsumerRecord<String, String> r : records) {
        OrderEvent event = JSON.parseObject(r.value(), OrderEvent.class);
        esService.indexOrder(event);      // 同步写 ES,偶尔要几百毫秒
    }
}

出事前的改动:ES 集群那天早上在做 segment 合并,写入耗时从平均 12 ms 涨到 380 ms。max.poll.records 默认 500 条,一批要处理 500 × 380ms = 190 秒,远超 max.poll.interval.ms 默认的 300 秒……

等等,190 秒小于 300 秒,那不该超时。实际情况更糟:有 3 个 @KafkaListener 线程,但它们是同一个 KafkaConsumer 实例吗?不是。concurrency = "3" 会创建 3 个独立的 consumer 线程,各自 poll。所以单线程一批 190 秒,还没到 300。

真正的触发点是:ES 那边偶发的 3 到 8 秒慢写。一批 500 条里只要混进 100 条慢的,总耗时就超过 300 秒了。日志里那次超时的批次,平均每条 612 ms。

而一旦被踢出组,rebalance 需要重新分配分区、重新 join,这期间消费完全停止。更糟的是 rebalance 之后新的一批还是处理不完,又超时——我们陷入了 rebalance 死循环。从 7:41 到 8:12,31 分钟里 rebalance 了 47 次,几乎没消费任何消息。

紧急处理

三步,按顺序做的:

# 1. 先把消费者停掉,避免继续 rebalance 消耗资源
$ kubectl scale deployment search-index --replicas=0

# 2. 改配置,把单批处理量降下来、把 poll 间隔放宽
# 3. 重新启动,并且临时增加实例数追赶积压
$ kubectl scale deployment search-index --replicas=12

改的配置:

spring:
  kafka:
    consumer:
      max-poll-records: 100          # 从默认 500 降到 100
      # 下面这两个不在 spring.kafka.consumer 的默认 key 里,用 properties 前缀
    properties:
      max.poll.interval.ms: 600000   # 5 分钟 → 10 分钟
      session.timeout.ms: 25000
      heartbeat.interval.ms: 8000
    listener:
      ack-mode: MANUAL               # 手动提交,处理完再提交

重启之后 12 个实例一起消费,lag 从 320 万降下来的过程:

时间     lag       消费速率
08:35    3,201,477   0(刚启动,在 join group)
08:40    2,988,230   42,500 条/分
09:00    2,105,884   44,100 条/分
09:30      802,331   43,500 条/分
10:05        8,204   43,800 条/分
10:12            0   -

1 小时 37 分钟追平。平均 43,700 条/分钟,也就是 728 条/秒。

分区数和消费者数的关系,这个必须搞懂

这是 Kafka 消费并行度的核心规则:

一个分区在同一时刻只能被消费者组里的一个消费者消费。所以分区数决定了消费者组的最大并行度

分区数消费者实例数实际分配结果
123每个消费者 4 个分区正常
1212每个消费者 1 个分区并行度打满
121612 个各拿 1 个,4 个空闲多出来的 4 个白跑
1211 个消费者消费全部 12 个分区串行

我们当时是 12 个分区、4 个实例、每个实例 concurrency=3,也就是 12 个消费者线程,刚好打满。所以临时扩容到 12 个实例时,每个实例只有 1 个线程有分区,剩下 2 个线程空闲。这是扩容效果不理想的原因之一。

真正的提速来自另一处:单个消费者的处理能力。降了 max-poll-records 之后反而不容易超时,加上 ES 恢复正常,单线程从 380 ms/条 降到 12 ms/条,这才是主要因素。

要扩分区吗

追平之后我评估过加分区。加分区的命令很简单:

$ bin/kafka-topics.sh --bootstrap-server 10.0.0.41:9092 \
    --alter --topic ORDER_EVENT_TOPIC --partitions 24

但有两个后果必须知道:

  • 改变 key 的分区映射。Kafka 用 hash(key) % partitionCount 决定消息去哪个分区。分区数变了,同一个 key 的消息会落到不同分区,分区内的顺序性在一段时间内会被破坏(老数据还在 12 个区,新数据去 24 个区)。我们这个 topic 用 orderId 做 key,依赖同订单消息的顺序,所以不能直接加。
  • 加分区不会重新分配已有数据,老分区的数据还在老分区。短期内会有数据倾斜。

我们最后的做法是新建一个 24 分区的 topic,生产者双写一段时间,等老的 topic 消费完再全切过去。麻烦但安全。

如果确定要原地加分区,务必先确认你的业务不依赖 key 的顺序性,或者消费端能容忍乱序。

批量消费调优

追平之后我把消费者的处理逻辑也改了。原来的写法是一条一条走网络,改成批量:

@KafkaListener(topics = "ORDER_EVENT_TOPIC", groupId = "search-index-group",
               concurrency = "3")
public void consume(List<ConsumerRecord<String, String>> records,
                    Acknowledgment ack) {
    try {
        List<OrderEvent> events = records.stream()
                .map(r -> JSON.parseObject(r.value(), OrderEvent.class))
                .collect(Collectors.toList());

        // 批量写 ES,一次网络往返
        esService.bulkIndex(events);

        ack.acknowledge();          // 处理成功再提交
    } catch (Exception e) {
        log.error("批量处理失败, size={}", records.size(), e);
        throw e;                     // 不提交,让 Kafka 重投
    }
}

批量前后:

方式单条耗时吞吐
逐条 index12 ms83 条/秒(单线程)
bulk(100 条)0.29 ms/条3450 条/秒(单线程)

41 倍。这个数字我在 ES 那篇里也提过,批量写是最有效的一招。

还有几个拉取端的参数:

spring:
  kafka:
    consumer:
      fetch-min-size: 102400        # 至少攒够 100 KB 才返回,减少空拉取
      fetch-max-wait: 500           # 最多等 500 ms
      max-poll-records: 100
    properties:
      max.poll.interval.ms: 600000

fetch.min.bytesfetch.max.wait.ms 是一对,满足任一条件就返回。默认是 1 字节、500 ms,意味着只要有 1 字节数据就立刻返回,网络往返很频繁。调成 100 KB 之后,网络请求数降了约 60%(我们单条消息平均 1.2 KB)。

慢消费者的另一种解法:poll + 处理线程池

如果单条消息的处理真的很慢(比如要调外部接口),把处理丢到线程池是常见做法。但这样会打乱 offset 提交的语义——线程还没处理完,offset 就提交了,消息可能丢。

正确做法是按分区维护处理队列,每个分区一个单线程池,保证分区内顺序,offset 按该分区实际完成的进度提交。Spring Kafka 2.4 里可以这么配:

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, String> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setConcurrency(3);
    factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
    return factory;
}

真要引入处理线程池的话,我倾向于另一条路:把慢逻辑拆出去,消费者只做"接收 + 落表 + 发一个内部 topic",慢处理交给下游。这样 offset 语义简单,也更好排查。我们最后选的是这条。

监控和告警

这次事故暴露的最大问题是没有 lag 告警。我们是靠业务方反馈才知道出事的,晚了半小时。

补上的监控:

# 每 30 秒采集一次所有消费者组的 lag,推给 Prometheus
$ bin/kafka-consumer-groups.sh --bootstrap-server 10.0.0.41:9092 \
    --describe --all-groups | awk 'NR>2 {sum += $6} END {print sum}'

告警规则:

  • lag 绝对值 > 50 万 → 严重告警,电话
  • lag 持续增长超过 10 分钟 → 警告,钉钉
  • 消费者组成员数为 0 → 严重告警(这次就是这种情况,最该早发现)
  • 5 分钟内 rebalance 次数 > 3 → 警告

最后一条是我们事后加的。rebalance 频繁几乎是"消费端出问题"的最早信号,比 lag 涨起来更早。

留个问题

关于《Kafka 消费积压的排查与扩容处理》里这个坑,你当时是怎么处理的?欢迎在评论区聊聊你踩过的类似情况。

参考