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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 05:40:08