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环境复现。
核心疑问
- 是否存在Quarkus批量确认消息导致超时的可能?
- 还有哪些其他排查方向?
分析与解答
关于批量确认导致超时的可能性
默认情况下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
相关产品推荐
相关产品推荐

