如何用Spring Kafka实现带重试、DLQ的简单Java生产者消费者程序
Spring Kafka 集成示例(含重试+死信队列)
1. 核心依赖(Maven)
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> <version>2.7.14</version> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> <version>2.9.10</version> </dependency> </dependencies>
2. application.yml 配置
spring: kafka: bootstrap-servers: 127.0.0.1:9092 # 生产者配置 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer acks: 1 retries: 3 # 消费者配置 consumer: group-id: business-group key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: earliest enable-auto-commit: false # 业务主题与死信主题定义 topic: business: business-topic dlq: business-topic-dlq
3. Kafka 核心配置类
这里配置重试策略和死信队列路由规则,重试次数设为3次,重试失败后自动将消息投递到死信队列:
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.*; import org.springframework.kafka.listener.DeadLetterPublishingRecoverer; import org.springframework.kafka.listener.DefaultErrorHandler; import org.springframework.kafka.support.ExponentialBackOffWithMaxRetries; import org.springframework.kafka.support.serializer.ErrorHandlingDeserializer; import java.util.HashMap; import java.util.Map; @Configuration public class KafkaConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServers; @Value("${spring.kafka.consumer.group-id}") private String groupId; @Value("${spring.kafka.topic.dlq}") private String dlqTopic; // 生产者工厂配置 @Bean public ProducerFactory<String, String> producerFactory() { Map<String, Object> configs = new HashMap<>(); configs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); configs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return new DefaultKafkaProducerFactory<>(configs); } @Bean public KafkaTemplate<String, String> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } // 消费者工厂配置 @Bean public ConsumerFactory<String, String> consumerFactory() { Map<String, Object> configs = new HashMap<>(); configs.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); configs.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); configs.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // 配置异常反序列化处理器,避免消费到非法格式消息直接宕机 configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); configs.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class); configs.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, StringDeserializer.class); return new DefaultKafkaConsumerFactory<>(configs); } // 监听器容器工厂,配置重试与死信队列 @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setConcurrency(3); // 配置退避策略:重试3次,每次间隔1秒 ExponentialBackOffWithMaxRetries backOff = new ExponentialBackOffWithMaxRetries(3); backOff.setInitialInterval(1000L); backOff.setMultiplier(1.0); // 重试失败后将消息投递到死信队列 DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate(), (consumerRecord, e) -> new TopicPartition(dlqTopic, consumerRecord.partition())); // 配置异常处理器 DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer, backOff); // 可指定需要重试的异常类型,不需要重试的异常直接进DLQ errorHandler.addRetryableExceptions(RuntimeException.class); errorHandler.addNotRetryableExceptions(IllegalArgumentException.class); factory.setCommonErrorHandler(errorHandler); return factory; } }
4. 生产者实现
import org.springframework.beans.factory.annotation.Value; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Component; @Component public class KafkaProducerService { private final KafkaTemplate<String, String> kafkaTemplate; @Value("${spring.kafka.topic.business}") private String businessTopic; public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendBusinessMessage(String message) { kafkaTemplate.send(businessTopic, message); } }
5. 消费者实现
5.1 业务消息消费者
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; @Component public class BusinessConsumer { @KafkaListener(topics = "${spring.kafka.topic.business}", containerFactory = "kafkaListenerContainerFactory") public void consumeBusinessMessage(ConsumerRecord<String, String> record) { String message = record.value(); System.out.println("收到业务消息:" + message); // 模拟消费异常,触发重试 if (message.contains("error")) { throw new RuntimeException("业务处理失败,触发重试"); } System.out.println("业务消息处理成功:" + message); } }
5.2 死信队列消费者
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; @Component public class DlqConsumer { @KafkaListener(topics = "${spring.kafka.topic.dlq}", groupId = "dlq-group") public void consumeDlqMessage(ConsumerRecord<String, String> record) { System.out.println("收到死信消息,key:" + record.key() + ",value:" + record.value() + ",可在此处做兜底处理或告警"); // 此处可以实现死信消息的人工干预触发、告警通知、持久化存储等逻辑 } }
6. 测试接口
import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; @RestController public class TestController { private final KafkaProducerService producerService; public TestController(KafkaProducerService producerService) { this.producerService = producerService; } @GetMapping("/send") public String sendMessage(@RequestParam String content) { producerService.sendBusinessMessage(content); return "消息发送成功"; } }
测试说明
- 本地启动Kafka服务,提前创建
business-topic和business-topic-dlq两个主题 - 启动Spring Boot服务,调用
http://localhost:8080/send?content=test可看到正常消费日志 - 调用
http://localhost:8080/send?content=test_error可看到连续3次重试日志,之后死信队列消费者收到异常消息
内容的提问来源于stack exchange,提问作者James Bond
相关产品推荐
相关产品推荐

