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

Spring Kafka批量消费下替代@RetryableTopic实现失败消息转DLT

Spring Kafka批量消费实现重试+DLT替代@RetryableTopic

由于@RetryableTopic暂不支持批量监听,我们可以通过手动结合RetryTemplate实现重试逻辑+KafkaTemplate发送失败消息到DLT的方式,复刻原有@RetryableTopic的核心能力:指定异常重试、指数退避、重试失败转DLT。

1. 配置RetryTemplate(对应原@RetryableTopic的重试规则)

创建配置类,定义符合需求的RetryTemplate,包括重试次数、退避策略、指定重试的异常类型:

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.retry.RetryPolicy;
import org.springframework.retry.backoff.ExponentialBackOffPolicy;
import org.springframework.retry.policy.SimpleRetryPolicy;
import org.springframework.retry.support.RetryTemplate;

import java.util.HashMap;
import java.util.Map;

@Configuration
public class KafkaRetryConfig {

    @Bean
    public RetryTemplate kafkaRetryTemplate() {
        RetryTemplate retryTemplate = new RetryTemplate();

        // 配置重试策略:仅对指定异常重试,重试4次
        Map<Class<? extends Throwable>, Boolean> retryableExceptions = new HashMap<>();
        retryableExceptions.put(ResourceAccessException.class, true);
        retryableExceptions.put(MyCustomRetryableException.class, true);
        RetryPolicy retryPolicy = new SimpleRetryPolicy(4, retryableExceptions);
        retryTemplate.setRetryPolicy(retryPolicy);

        // 配置退避策略:初始延迟1000ms,乘数2.0(指数退避)
        ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy();
        backOffPolicy.setInitialInterval(1000);
        backOffPolicy.setMultiplier(2.0);
        retryTemplate.setBackOffPolicy(backOffPolicy);

        return retryTemplate;
    }
}

2. 批量消费者实现(重试+DLT转发)

在批量监听方法中,遍历每条消息,用RetryTemplate执行消费逻辑,重试失败后将消息发送到DLT主题:

import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.retry.support.RetryTemplate;
import org.springframework.stereotype.Component;

import java.util.List;

@Component
public class BatchKafkaConsumer {

    private final RetryTemplate kafkaRetryTemplate;
    private final KafkaTemplate<String, Object> kafkaTemplate;
    // DLT主题名称,需手动创建
    private static final String DLT_TOPIC = "original-topic-dlt";

    public BatchKafkaConsumer(RetryTemplate kafkaRetryTemplate, KafkaTemplate<String, Object> kafkaTemplate) {
        this.kafkaRetryTemplate = kafkaRetryTemplate;
        this.kafkaTemplate = kafkaTemplate;
    }

    @KafkaListener(topics = "original-topic", batch = "true", containerFactory = "batchKafkaListenerContainerFactory")
    public void batchConsume(List<Object> messages) {
        for (Object message : messages) {
            try {
                // 用RetryTemplate执行消费逻辑,自动处理重试
                kafkaRetryTemplate.execute(context -> {
                    processMessage(message);
                    return null;
                });
            } catch (Exception e) {
                // 重试耗尽后,将消息发送到DLT
                sendToDlt(message, e);
            }
        }
    }

    /**
     * 实际的消息处理逻辑
     */
    private void processMessage(Object message) {
        // 此处编写业务逻辑,可能抛出ResourceAccessException或MyCustomRetryableException
        // 示例:throw new MyCustomRetryableException("处理失败");
    }

    /**
     * 将失败消息发送到DLT
     */
    private void sendToDlt(Object message, Exception cause) {
        // 可添加失败原因到消息头/消息体,方便后续排查
        kafkaTemplate.send(DLT_TOPIC, message);
    }
}

3. 批量消费容器配置

确保批量消费容器正确开启批量监听:

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;

@Configuration
public class KafkaBatchConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> batchKafkaListenerContainerFactory(ConsumerFactory<String, Object> consumerFactory) {
        ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        // 开启批量监听
        factory.setBatchListener(true);
        // 设置拉取超时时间
        factory.getContainerProperties().setPollTimeout(3000);
        return factory;
    }
}

4. application.yml配置示例

spring:
  kafka:
    consumer:
      bootstrap-servers: your-kafka-server:9092
      group-id: your-consumer-group
      max-poll-records: 5
      auto-offset-reset: earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "*"
    producer:
      bootstrap-servers: your-kafka-server:9092
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer

注意事项

  • 需手动创建DLT主题original-topic-dlt(因autoCreateTopics=false)
  • 若需对批量消息整体重试而非单个,可调整逻辑将整个batch传入RetryTemplate,但注意批量重试可能导致已成功消息重复处理,仅在业务允许时使用
  • 可在sendToDlt方法中添加异常信息到消息头,便于后续问题排查

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 10:10:30