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
相关产品推荐
相关产品推荐

