You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Quarkus处理Kafka消息时遇TooManyMessagesWithoutAckException问题排查

Quarkus Kafka消费超时问题排查

问题背景

我们在Quarkus进程中消费Kafka消息时执行以下步骤:

  • Thread.sleep(30000) - 业务逻辑要求
  • 调用第三方API
  • 调用另一个第三方API
  • 向数据库插入数据

该进程几乎每日都会抛出TooManyMessagesWithoutAckException后挂起,相关日志如下:

2022-12-02 20:02:50 INFO  [2bdf7fc8-e0ad-4bcb-87b8-c577eb506b38,     ] : Going to sleep for 30 sec.....
2022-12-02 20:03:20 WARN  [                    kafka] : SRMSG18231: The record 17632 from topic-partition '<partition>' has waited for 60 seconds to be acknowledged. This waiting time is greater than the configured threshold (60000 ms). At the moment 2 messages from this partition are awaiting acknowledgement. The last committed offset for this partition was 17631. This error is due to a potential issue in the application which does not acknowledged the records in a timely fashion. The connector cannot commit as a record processing has not completed.
2022-12-02 20:03:20 WARN  [                     kafka] : SRMSG18228: A failure has been reported for Kafka topics '[<topic name>]': io.smallrye.reactive.messaging.kafka.commit.KafkaThrottledLatestProcessedCommit$TooManyMessagesWithoutAckException: The record 17632 from topic/partition '<partition>' has waited for 60 seconds to be acknowledged. At the moment 2 messages from this partition are awaiting acknowledgement. The last committed offset for this partition was 17631.
2022-12-02 20:03:20 INFO  [2bdf7fc8-e0ad-4bcb-87b8-c577eb506b38,     ] : Sleep over!

消息消费代码示例:

@Incoming("my-channel")
@Blocking
CompletionStage<Void> consume(Message<Person> person) {
     String msgKey = (String) person
        .getMetadata(IncomingKafkaRecordMetadata.class).get()
        .getKey();
        // ... 执行业务步骤:sleep、第三方API调用、DB插入
      return person.ack();
}

根据日志,消息拉取后仅耗时30秒,却触发了“未在60秒内确认消息”的异常。排查当日日志,未发现第三方API调用耗时超过30秒的情况。

当前Kafka配置如下:

mp:
  messaging:
    incoming:
      my-channel:
        topic: <topic>
        group:
          id: <group id>
        connector: smallrye-kafka
        value:
          serializer: org.apache.kafka.common.serialization.StringSerializer
          deserializer: org.apache.kafka.common.serialization.StringDeserializer

该Kafka主题有4个分区,副本因子为3,对应运行3个Quarkus进程Pod。此问题无法在Dev和UAT环境复现。

核心疑问

  1. 是否存在Quarkus批量确认消息导致超时的可能?
  2. 还有哪些其他排查方向?

分析与解答

关于批量确认导致超时的可能性

默认情况下SmallRye Reactive Messaging Kafka使用KafkaThrottledLatestProcessedCommit策略,该策略跟踪分区最新已处理偏移量并定期提交,但超时异常的核心是单条记录的等待确认时间超过阈值(60秒),和批量提交本身无关。不过存在一种关联场景:
如果同一个分区的多条消息被拉取后,前面的消息占用线程池资源处理,后续消息会排队等待处理,此时后续消息的等待时间从被拉取时开始计算,而非自身开始处理时。比如:

  • 消息A被拉取,开始处理(sleep30秒+其他操作)
  • 消息A处理过程中,消息B被拉取并进入线程池队列等待
  • 消息A处理完成耗时35秒,此时消息B已等待35秒,加上自身处理的30秒,总等待时间超过60秒触发异常

日志中提到的“At the moment 2 messages from this partition are awaiting acknowledgement”,正好匹配这种线程排队导致的时间累积场景。

其他排查方向

  • 线程池配置检查:@Blocking默认使用Quarkus全局线程池,若线程池大小不足,会导致同分区消息排队等待时间拉长。可以为消费者指定专用线程池:
    @Blocking("kafka-consumer-pool")
    
    并配置线程池参数:
    quarkus:
      thread-pool:
        kafka-consumer-pool:
          core-threads: 10
          max-threads: 20
    
  • Kafka消费者拉取配置:调小max.poll.records(默认500),减少单批次拉取的消息数量,避免线程池被瞬间占满。比如设置为10或20:
    mp:
      messaging:
        incoming:
          my-channel:
            max.poll.records: 10
    
  • 分区与Pod负载均衡:主题有4个分区但仅运行3个Pod,会有一个Pod分配到2个分区,负载不均易引发线程压力。调整Pod数量为4,让每个Pod处理一个分区,均衡负载。
  • 异常处理逻辑检查:确保代码中存在完善的异常捕获机制,异常时调用nack()方法,避免消息因未处理完成一直处于等待确认状态。
  • 生产环境特定因素排查:Dev/UAT无法复现可能是因为生产环境消息量更大、第三方API延迟波动累积、网络抖动等。添加详细监控:
    • 记录每条消息的拉取时间、开始处理时间、处理完成时间、确认时间,分析时间线
    • 监控线程池的活跃线程数、队列长度
    • 监控第三方API的响应时间分布(比如99分位延迟)
  • 提交策略参数调整:若业务处理确实需要较长时间,可临时调大max.wait.time.ms(默认60秒),但这只是缓解手段,核心仍需解决消息排队问题:
    mp:
      messaging:
        incoming:
          my-channel:
            max.wait.time.ms: 120000
    

内容的提问来源于stack exchange,提问作者KVS

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 03:05:42