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

如何用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 "消息发送成功";
    }
}

测试说明

  1. 本地启动Kafka服务,提前创建business-topic和business-topic-dlq两个主题
  2. 启动Spring Boot服务,调用http://localhost:8080/send?content=test可看到正常消费日志
  3. 调用http://localhost:8080/send?content=test_error可看到连续3次重试日志,之后死信队列消费者收到异常消息

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 00:51:00