Spring Kafka 3.3响应式@KafkaListener异常时偏移量立即提交问题
问题描述
使用Spring Kafka 3.3.10搭配返回Mono的响应式Kafka监听器,监听器调用下游WebClient并返回Mono。发生超时异常时,期望不提交偏移量,执行重试或发送至DLQ,但实际异常发生后偏移量立即提交,DefaultErrorHandler未触发重试。
监听器代码
@KafkaListener( topics = "internal.request.topic", groupId = "reactive-demo", containerFactory = "kafkaListenerContainerFactory" ) public Mono<Void> consume(String message) { return downstreamClient.call(message) .timeout(Duration.ofSeconds(3)) .onErrorMap(TimeoutException.class, e -> new DownstreamTimeoutException("Timeout")) .then(); }
监听器容器配置
@Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory( DefaultErrorHandler handler) { var factory = new ConcurrentKafkaListenerContainerFactory<String, String>(); factory.setConcurrency(3); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD); factory.setCommonErrorHandler(handler); return factory; }
错误处理器配置
@Bean public CommonErrorHandler handler() { DefaultErrorHandler errorHandler = new DefaultErrorHandler(new FixedBackOff(2000L, 2L)); errorHandler.addRetryableExceptions(DownstreamTimeoutException.class, TimeoutException.class, ListenerExecutionFailedException.class); errorHandler.setRetryListeners((record, ex, attempt) -> log.warn("[CUSTOM-RETRY] attempt={} topic={} partition={} offset={} ex={}", attempt, record.topic(), record.partition(), record.offset(), ex.toString()) ); return errorHandler; }
观测到的行为
- 超时异常发生时,监听器记录错误日志,但偏移量立即提交
- DefaultErrorHandler的重试逻辑未执行,自定义重试日志无输出
- 消息直接被确认,未进入重试流程
预期行为
- 针对
DownstreamTimeoutException:不提交偏移量,按配置重试,重试耗尽后发送至DLQ - 其他异常:处理完成后提交偏移量
- 成功处理时正常提交偏移量
原因分析
当前使用的ConcurrentKafkaListenerContainerFactory是为阻塞式监听器设计的,在响应式场景下,它会默认在**Mono完成(无论成功还是失败)**后立即提交偏移量,导致DefaultErrorHandler无法拦截异常并控制偏移量提交逻辑。此外,ListenerExecutionFailedException是Spring Kafka包装的异常,默认会被判定为非重试异常,进一步阻碍重试流程触发。
解决方案
1. 切换为响应式容器工厂
替换ConcurrentKafkaListenerContainerFactory为ReactiveKafkaListenerContainerFactory,这是Spring Kafka专为响应式监听器设计的容器,能正确关联DefaultErrorHandler并控制偏移量提交:
@Bean public ReactiveKafkaListenerContainerFactory<String, String> reactiveKafkaListenerContainerFactory( ConsumerFactory<String, String> consumerFactory, DefaultErrorHandler errorHandler) { ReactiveKafkaListenerContainerFactory<String, String> factory = new ReactiveKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setBatch(false); // 单条消费模式 factory.setCommonErrorHandler(errorHandler); // 配置确认模式 factory.getContainerOptions().setAckMode(ContainerProperties.AckMode.RECORD); return factory; }
同时更新监听器的容器工厂引用:
@KafkaListener( topics = "internal.request.topic", groupId = "reactive-demo", containerFactory = "reactiveKafkaListenerContainerFactory" // 切换为响应式容器工厂 ) public Mono<Void> consume(String message) { return downstreamClient.call(message) .timeout(Duration.ofSeconds(3)) .onErrorMap(TimeoutException.class, e -> new DownstreamTimeoutException("Timeout")) .then(); }
2. 调整错误处理器配置
移除对ListenerExecutionFailedException的重试配置,直接针对实际业务异常设置重试逻辑,并添加DLQ发送器:
@Bean public CommonErrorHandler errorHandler(KafkaOperations<String, String> kafkaOperations) { // 配置DLQ发送器,重试耗尽后发送到指定主题 DeadLetterPublishingRecoverer dlqRecoverer = new DeadLetterPublishingRecoverer(kafkaOperations, (record, ex) -> new TopicPartition("internal.request.topic.dlq", record.partition())); DefaultErrorHandler errorHandler = new DefaultErrorHandler(dlqRecoverer, new FixedBackOff(2000L, 2L)); // 指定需要重试的异常类型 errorHandler.addRetryableExceptions(DownstreamTimeoutException.class, TimeoutException.class); // 可选:排除不需要重试的异常 // errorHandler.addNotRetryableExceptions(IllegalArgumentException.class); errorHandler.setRetryListeners((record, ex, attempt) -> log.warn("[CUSTOM-RETRY] attempt={} topic={} partition={} offset={} ex={}", attempt, record.topic(), record.partition(), record.offset(), ex.toString()) ); return errorHandler; }
3. 验证逻辑
- 响应式容器会在Mono成功完成时自动提交偏移量
- 当异常被
DefaultErrorHandler拦截时,会暂停偏移量提交,执行重试逻辑 - 重试耗尽后,消息会被发送至配置的DLQ主题
内容的提问来源于stack exchange,提问作者Sameep Karanjkar
相关产品推荐
相关产品推荐

