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

Reactor Kafka开启乱序提交后消费无限停滞问题求助

问题:Reactor-Kafka开启乱序提交后消费无限停滞

我们使用ReactiveKafkaConsumerTemplate接收消息,处理完成后手动确认offset。开启乱序提交(maxDeferredCommits=250)后,在特定场景下出现消费无限停滞的问题,事件流程如下:

  • 网络故障或Kafka服务器维护触发RetriableCommitFailedException
  • 消费者因“Paused - commits are retrying”暂停消费
  • 消费者恢复后开始“Emitting records”,但后续无“Async committing”日志(未触发新的提交操作)
  • 多次拉取消息后,消费者因“Paused - too many deferred commits”再次暂停
  • 此后无ConsumerEventLoop相关日志输出,直到发生Rebalance(3台主机部署3个消费者,移除1台即可恢复消费)

相关配置及版本信息

reactor-kafka-1.3.13.jar
logging:
 level:
   reactor:
     kafka:
       receiver: DEBUG
maxDeferredCommits: 250
ConsumerConfig
auto.commit.interval.ms = 1000
auto.offset.reset = earliest
connections.max.idle.ms = 540000
enable.auto.commit = false
heartbeat.interval.ms = 1000
max.poll.interval.ms = 300000     
max.poll.records = 500
request.timeout.ms = 30000
session.timeout.ms = 10000

关键日志片段

11/24/22 6:50:06.386 AM DEBUG r.k.r.internals.ConsumerEventLoop Async committing: {
test-0=OffsetAndMetadata{offset=12206778, leaderEpoch=null, metadata=''}, 
test-1=OffsetAndMetadata{offset=12253822, leaderEpoch=null, metadata=''} 
test-2=OffsetAndMetadata{offset=12257066, leaderEpoch=null, metadata=''} 
test-3=OffsetAndMetadata{offset=12265134, leaderEpoch=null, metadata=''}}

No more “Async committing” after this

11/24/22 6:50:06.451 AM WARN r.k.r.internals.ConsumerEventLoop Commit failed with org.apache.kafka.clients.consumer.RetriableCommitFailedException: Offset commit failed with a retriable exception. You should retry committing the latest consumed offsets. Caused by: org.apache.kafka.common.errors.DisconnectException: null

11/24/22 6:50:06.452 AM WARN r.k.r.internals.ConsumerEventLoop Commit failed with exceptionorg.apache.kafka.clients.consumer.RetriableCommitFailedException: Offset commit failed with a retriable exception. You should retry committing the latest consumed offsets., retries remaining 99
…
11/24/22 6:50:06.452 AM WARN r.k.r.internals.Commit failed with exceptionorg.apache.kafka.clients.consumer.RetriableCommitFailedException: Offset commit failed with a retriable exception. You should retry committing the latest consumed offsets., retries remaining 93

11/24/22 6:50:06.486 DEBUG r.k.r.internals.ConsumerEventLoop -Paused - commits are retrying
11/24/22 6:50:06.987 DEBUG r.k.r.internals.ConsumerEventLoop -Resumed
11/24/22 6:50:07.387 DEBUG r.k.r.internals.ConsumerEventLoop -Emitting 1 records, requested now 1
11/24/22 6:50:07.387 DEBUG r.k.r.internals.ConsumerEventLoop -onRequest.toAdd 1, paused false
…

11/24/22 6:51:05.248 DEBUG r.k.r.internals.ConsumerEventLoop -Paused - too many deferred commits
11/24/22 6:51:05.248 DEBUG r.k.r.internals.ConsumerEventLoop -Consumer woken

No more “ConsumerEventLoop” log after this until rebalance

核心消费代码

consumeMessge() {

ReceiverOptions basicReceiverOptions = ReceiverOptions.create(
   consumerProperties)
   .maxDeferredCommits(250)
   .commitInterval(Duration.ofMillis(commitInterval))
   .subscription(topics);

reactiveKafkaConsumerTemplate=new ReactiveKafkaConsumerTemplate<>(basicReceiverOptions);

 return reactiveKafkaConsumerTemplate
     .receive()
     .publishOn(Schedulers.boundedElastic())
     .flatMap(x -> Mono.just(x)
              .delayElement(Duration.ofMillis(500),10))
     .flatMap(receiverRecord ->
            //process the record
             messageServiceImpl.process(receiverRecord)
             .doFinally(x -> {
               //ack offset
               log.info("MessageConsumer ACK offset={} ", receiverRecord.offset());
               receiverRecord.receiverOffset().acknowledge();
             })
             .subscribeOn(Schedulers.boundedElastic())
     )
     .....

}

原因分析

  1. 提交重试阻塞提交队列:当RetriableCommitFailedException发生后,消费者进入提交重试逻辑,此时消费暂停;恢复后继续拉取消息并ACK,但之前失败的提交任务可能阻塞了内部的提交队列,导致新的ACK无法触发异步提交操作。
  2. Deferred Commits累积超限:由于提交队列阻塞,已ACK的offset无法被批量提交,deferred commits数量不断累积,达到maxDeferredCommits=250阈值后,消费者再次暂停消费。此时因为没有新的提交触发逻辑,消费彻底停滞。
  3. 版本潜在bug:reactor-kafka 1.3.13版本在处理提交重试与乱序提交的交互逻辑上可能存在缺陷,导致提交队列无法正常恢复,进而引发消费停滞。
  4. 线程调度问题:代码中使用publishOn和subscribeOn切换线程,可能导致ACK操作的上下文与消费者事件循环的上下文不一致,影响提交信号的传递。

解决方案

  1. 升级reactor-kafka版本:升级至最新稳定版(如1.4.x或更高),该版本已修复多个提交相关的bug,优化了提交重试与乱序提交的交互逻辑。
  2. 调整乱序提交参数:
    • 适当调大maxDeferredCommits值,避免因短暂提交阻塞导致快速触发暂停阈值;
    • 缩短commitInterval,增加提交频率,减少deferred commits的累积速度。
  3. 优化ACK逻辑:确保ACK操作在消息处理完成的正确时机执行,避免在doFinally中可能出现的线程上下文问题,可尝试将ACK操作移至process方法的成功回调中,或使用then()确保处理完成后再ACK:
    .flatMap(receiverRecord -> 
        messageServiceImpl.process(receiverRecord)
            .then(Mono.fromRunnable(() -> {
                log.info("MessageConsumer ACK offset={} ", receiverRecord.offset());
                receiverRecord.receiverOffset().acknowledge();
            }))
            .subscribeOn(Schedulers.boundedElastic())
    )
    
  4. 配置提交重试策略:在ReceiverOptions中配置提交重试的间隔和最大重试次数,避免短时间内大量重试导致队列阻塞:
    .commitRetryInterval(Duration.ofSeconds(1))
    .maxCommitRetries(10)
    
  5. 增加Rebalance触发机制:可通过监控消费者状态,当检测到长时间无消费日志时,主动触发Rebalance(如重启消费者实例),作为临时恢复方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 03:40:34