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()) ) ..... }
原因分析
- 提交重试阻塞提交队列:当
RetriableCommitFailedException发生后,消费者进入提交重试逻辑,此时消费暂停;恢复后继续拉取消息并ACK,但之前失败的提交任务可能阻塞了内部的提交队列,导致新的ACK无法触发异步提交操作。 - Deferred Commits累积超限:由于提交队列阻塞,已ACK的offset无法被批量提交,
deferred commits数量不断累积,达到maxDeferredCommits=250阈值后,消费者再次暂停消费。此时因为没有新的提交触发逻辑,消费彻底停滞。 - 版本潜在bug:reactor-kafka 1.3.13版本在处理提交重试与乱序提交的交互逻辑上可能存在缺陷,导致提交队列无法正常恢复,进而引发消费停滞。
- 线程调度问题:代码中使用
publishOn和subscribeOn切换线程,可能导致ACK操作的上下文与消费者事件循环的上下文不一致,影响提交信号的传递。
解决方案
- 升级reactor-kafka版本:升级至最新稳定版(如1.4.x或更高),该版本已修复多个提交相关的bug,优化了提交重试与乱序提交的交互逻辑。
- 调整乱序提交参数:
- 适当调大
maxDeferredCommits值,避免因短暂提交阻塞导致快速触发暂停阈值; - 缩短
commitInterval,增加提交频率,减少deferred commits的累积速度。
- 适当调大
- 优化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()) ) - 配置提交重试策略:在
ReceiverOptions中配置提交重试的间隔和最大重试次数,避免短时间内大量重试导致队列阻塞:.commitRetryInterval(Duration.ofSeconds(1)) .maxCommitRetries(10) - 增加Rebalance触发机制:可通过监控消费者状态,当检测到长时间无消费日志时,主动触发Rebalance(如重启消费者实例),作为临时恢复方案。
内容的提问来源于stack exchange,提问作者Gin
相关产品推荐
相关产品推荐

