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

Spring Kafka 3.0.6阻塞重试实现代码检查:是否存在遗漏?

Spring Kafka 3.0.6 阻塞重试与死信队列代码检查及修正

当前代码目标是实现:处理消息遇到RetryableException时重试2次,重试耗尽后将消息发送至指定固定分区的死信队列。现有代码存在以下遗漏和问题:

  • 异常范围未限定:DefaultErrorHandler默认会对除FatalThrowable外的所有异常重试,不符合仅对RetryableException重试的需求。
  • 重试次数配置错误:FixedBackOff(0L, 2L)的maxAttempts参数代表总尝试次数(含第一次执行),设置为2L仅会触发1次重试,无法满足“重试2次”的要求。
  • Header参数类型不匹配:消费方法中partition、offset、timestamp定义为String,但Kafka实际传递的是Integer(partition)、Long(offset、timestamp),会触发类型转换异常。
  • 死信队列Topic硬编码:死信队列名称写死为"topicName" + ".DLT",与消费的demoTopicName不对应,易引发错误。
  • 非RetryableException处理遗漏:捕获非重试异常后仅打印日志,未提交偏移量,会导致消息重复消费。

修正后的消费方法代码

@KafkaListener(
        autoStartup = "false",
        containerFactory = "concurrentKafkaListenerContainerFactory",
        id = "demoConsumer",
        groupId = "demoConsumerGroup",
        topics = "demoTopicName"
)
public void consumeMessage(
        @Header(KafkaHeaders.RECEIVED_PARTITION) Integer partition,
        @Header(KafkaHeaders.OFFSET) Long offset,
        @Header(KafkaHeaders.RECEIVED_TIMESTAMP) Long timestamp,
        ConsumerRecord<String, String> consumerRecord,
        Acknowledgment acknowledgment) {

    try {
        log.info("Processing message from partition: {}, offset: {}", partition, offset);
        // 业务处理逻辑
        acknowledgment.acknowledge();
    } catch (Exception ex) {
        if (ex instanceof RetryableException) {
            throw ex; // 抛出异常触发重试逻辑
        }
        log.error("Non-retryable exception occurred", ex);
        acknowledgment.acknowledge(); // 提交偏移量,避免重复消费
    }
}

修正后的配置类代码

@EnableKafka
@Configuration
@RequiredArgsConstructor
public class KafkaConfig {

    private final KafkaTemplate<String, String> kafkaTemplate;

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String>
    concurrentKafkaListenerContainerFactory(ConsumerFactory<String, String> consumerFactory) {

        ConcurrentKafkaListenerContainerFactory<String, String> factory =
                new ConcurrentKafkaListenerContainerFactory<>();

        factory.setConsumerFactory(consumerFactory);
        factory.setConcurrency(1);
        factory.getContainerProperties().setStopImmediate(true);
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
        factory.setCommonErrorHandler(defaultErrorHandler());

        return factory;
    }

    @Bean
    public CommonErrorHandler defaultErrorHandler() {
        // 总尝试次数3次 = 1次初始执行 + 2次重试
        DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer(), new FixedBackOff(0L, 3L));
        // 仅对RetryableException触发重试
        errorHandler.addRetryableExceptions(RetryableException.class);
        // 排除所有其他异常,不进行重试
        errorHandler.addNotRetryableExceptions(Exception.class);
        return errorHandler;
    }

    @Bean
    public DeadLetterPublishingRecoverer recoverer() {
        final BiFunction<ConsumerRecord<?, ?>, Exception, TopicPartition>
                CUSTOMIZE_DESTINATION_RESOLVER = (cr, e) ->
                        // 根据原Topic动态生成死信队列,固定发送至分区0
                        new TopicPartition(cr.topic() + ".DLT", 0);

        return new DeadLetterPublishingRecoverer(kafkaTemplate, CUSTOMIZE_DESTINATION_RESOLVER);
    }
}

关键修正说明

  1. 异常精确控制:通过addRetryableExceptions和addNotRetryableExceptions限定仅RetryableException触发重试。
  2. 重试次数校准:将FixedBackOff的maxAttempts设为3L,确保初始执行失败后进行2次重试。
  3. 参数类型修正:将Header参数改为Kafka实际传递的类型,避免类型转换错误。
  4. 死信队列动态生成:根据原消息Topic自动生成死信队列名称,避免硬编码错误。
  5. 偏移量提交:非重试异常处理时提交偏移量,防止消息重复消费。
  6. Bean管理:将错误处理器和死信恢复器声明为Spring Bean,确保上下文正确管理。

内容的提问来源于stack exchange,提问作者Sk Monjurul Haque

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 21:07:43