文章目录

第一部分:基础概念

1.1 什么是 RabbitMQ

RabbitMQ 是一个开源的消息代理(Message Broker)队列服务器,实现了高级消息队列协议(AMQP 0.9.1)。

  • 核心价值:
    • 解耦: 系统 A 发送消息后无需等待系统 B 处理,降低耦合度。
    • 异步: 耗时操作(如发邮件、生成报表)放入队列,主流程立即返回,提升响应速度。
    • 削峰填谷: 突发流量先存入队列,消费者按能力慢慢处理,防止系统崩溃。

1.2 核心术语对照表

术语 英文 解释 生活类比
Broker Broker RabbitMQ 服务端实体 邮局大楼
Virtual Host VHost 虚拟主机,用于逻辑隔离 邮局里的不同分区(如国内件区、国际件区)
Exchange Exchange 交换机,接收消息并路由 分拣中心
Queue Queue 队列,存储消息的缓冲区 具体的邮箱/包裹架
Binding Binding 连接交换机和队列的规则 分拣规则(如:北京的去 A 区)
Routing Key Routing Key 路由键,消息携带的标签 包裹上的地址标签
Producer Producer 消息生产者 寄件人
Consumer Consumer 消息消费者 收件人

1.3 核心交互流程图

RabbitMQ Server (Broker)

1.发送消息 + 路由键

2.根据规则路由

3.推送/拉取

生产者 Producer

交换机 Exchange

队列 Queue

消费者 Consumer


第二部分:项目搭建与基础配置

2.1 标准项目结构

基于您的实战项目,采用清晰的分层结构:

spring-rabbitmq/
├── src/main/java/org/example/rabbitmq/
│   ├── config/                 # 【核心】所有 MQ 组件定义 (Exchange, Queue, Binding)
│   │   ├── DeadLetterRabbitConfig.java   # 死信队列配置 (异常兜底)
│   │   ├── DelayRabbitConfig.java        # 延迟队列配置 (TTL+DLX)
│   │   ├── DirectRabbitConfig.java       # 直连交换机 (点对点)
│   │   ├── FanoutRabbitConfig.java       # 扇形交换机 (广播)
│   │   ├── TopicRabbitConfig.java        # 主题交换机 (模糊匹配)
│   │   ├── MessageListenerConfig.java    # 手动确认监听器容器配置
│   │   └── RabbitConfig.java             # 全局 RabbitTemplate 及回调配置
│   ├── constant/               # 常量定义,避免硬编码
│   │   └── RabbitMqConstant.java
│   ├── rabbit/                 # 注解式接收者 (@RabbitListener)
│   │   ├── DelayReceiver.java             # 延迟消息处理
│   │   ├── DirectReceiver.java            # 直连消息处理
│   │   ├── FanoutReceiverA/B/C.java       # 广播消息处理
│   │   └── TopicMan/TotalReceiver.java    # 主题消息处理
│   ├── receiver/               # 接口式接收者 (ChannelAwareMessageListener)
│   │   └── MyAckReceiver.java             # 统一手动 ACK 逻辑
│   └── SpringRabbitmqStartApplication.java
├── src/main/resources/
│   ├── application.yaml        # 连接与基础行为配置
│   └── logback-spring.xml      # 日志配置
└── pom.xml

2.2 依赖配置 (pom.xml)

确保引入以下核心依赖,特别注意 JSON 序列化支持

<dependencies>
    <!-- Spring Boot RabbitMQ Starter -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
    <!-- 测试支持 -->
    <dependency>
        <groupId>org.springframework.amqp</groupId>
        <artifactId>spring-rabbit-test</artifactId>
        <scope>test</scope>
    </dependency>
    <!-- JSON 序列化支持 (生产环境必须,避免 JDK 默认序列化的坑) -->
    <dependency>
        <groupId>com.fasterxml.jackson.core</groupId>
        <artifactId>jackson-databind</artifactId>
    </dependency>
    <!-- 日志增强 (可选) -->
    <dependency>
        <groupId>net.logstash.logback</groupId>
        <artifactId>logstash-logback-encoder</artifactId>
    </dependency>
</dependencies>

2.3 核心配置文件 (application.yaml) 深度解析

spring:
  rabbitmq:
    host: 192.168.188.153
    port: 5672              # AMQP 协议端口 (管理界面是 15672)
    username: kiv
    password: kiv
    virtual-host: /rabbitmqkiv # 虚拟主机,实现多租户隔离
    
    # === 生产者可靠性保障 (防止消息丢失第一道防线) ===
    # correlated: 开启发布确认,消息到达交换机后回调 ConfirmCallback
    publisher-confirm-type: correlated 
    # true: 消息到达交换机但无法路由到队列时,触发 ReturnCallback
    publisher-returns: true 
    
    listener:
      simple:
        # === 消费者可靠性保障 (防止消息丢失第二道防线) ===
        # manual: 必须手动调用 channel.basicAck(),否则消息不会移除
        acknowledge-mode: manual 
        
        # 重试机制 (可选,建议结合死信队列使用)
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000ms

2.4 全局配置类 (RabbitConfig.java)

关键点:配置 JSON 序列化转换器,并注册确认回调。这是生产环境的标配。

@Configuration
public class RabbitConfig {

    @Bean
    public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
        RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
        
        // 1. 设置 JSON 序列化转换器 (重要!避免乱码和类型转换错误)
        rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter());

        // 2. 开启 ConfirmCallback (消息到达交换机)
        rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
            if (!ack) {
                System.err.println("❌ 消息发送失败,未到达交换机: " + 
                    (correlationData != null ? correlationData.getId() : "null") + ", 原因: " + cause);
                // TODO: 记录数据库,启动定时任务重发
            } else {
                System.out.println("✅ 消息成功到达交换机: " + 
                    (correlationData != null ? correlationData.getId() : "null"));
            }
        });

        // 3. 开启 ReturnsCallback (消息未路由到队列)
        rabbitTemplate.setReturnsCallback(returnedMessage -> {
            System.err.println("❌ 消息路由失败: " + new String(returnedMessage.getMessage().getBody()) 
                + ", 交换机: " + returnedMessage.getExchange() 
                + ", 路由键: " + returnedMessage.getRoutingKey());
            // TODO: 记录死信或报警
        });

        return rabbitTemplate;
    }
}

第三部分:五种交换机模式实战

3.1 直连交换机 (Direct Exchange)

  • 机制: 精确匹配Routing Key 必须与 Binding Key 完全一致。
  • 场景: 点对点通知、特定任务分发。
  • 注意: 如果 Routing Key 不匹配且未开启 ReturnCallback,消息会直接丢失。
测试类 (DirectSender.java)
@SpringBootTest
class DirectTests {

    @Autowired
    RabbitTemplate rabbitTemplate;  // 使用RabbitTemplate,这提供了接收/发送等等方法
    
    // 测试发送消息到Direct
    @Test
    public void sendDirectMessage() throws InterruptedException {
        // 循环发送消息
        while (true) {
            String messageId = String.valueOf(UUID.randomUUID());
            String messageData = "test message, hello!";
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String, Object> map = new HashMap<>();
            map.put("messageId", messageId);
            map.put("messageData", messageData);
            map.put("createTime", createTime);
            // 将消息携带绑定键值:TestDirectRouting 发送到交换机 TestDirectExchange
            rabbitTemplate.convertAndSend("TestDirectExchange", "TestDirectRouting", map);
            // sleep(50);
        }
    }
 

}
配置类 (DirectRabbitConfig.java)
@Configuration
public class DirectRabbitConfig {
    
    @Bean
    public Queue TestDirectQueue() {
        // durable=true: 队列持久化,RabbitMQ 重启后依然存在
        return new Queue(RabbitMqConstant.TEST_DIRECT_QUEUE, true);
    }

    @Bean
    DirectExchange TestDirectExchange() {
        return new DirectExchange(RabbitMqConstant.TEST_DIRECT_EXCHANGE, true, false);
    }

    @Bean
    Binding bindingDirect() {
        return BindingBuilder
                .bind(TestDirectQueue())
                .to(TestDirectExchange())
                .with(RabbitMqConstant.TEST_DIRECT_ROUTING); // 精确匹配
    }
}
消费者 (DirectReceiver.java)
@Component
public class DirectReceiver {
    Integer count = 0;

    @RabbitHandler
    public void process(Map testMessage) {
        count++;
        System.out.println("【Direct】消费者收到消息 [" + count + "]: " + testMessage.toString());
        // 注意:若全局配置了 manual ack,此处需配合 MessageListenerConfig 使用
    }
}

3.2 主题交换机 (Topic Exchange)

  • 机制: 模糊匹配。支持通配符:
    • * (星号): 匹配一个单词 (如 user.* 匹配 user.login)。
    • # (井号): 匹配零个或多个单词 (如 user.# 匹配 user.login.success)。
  • 场景: 复杂的日志收集、多级业务通知。
测试类 (TopicSender.java)
@SpringBootTest
class TopicTests {

    @Autowired
    RabbitTemplate rabbitTemplate;
    /**
     * queue: topic.man 和 topic.woman 都可以发送消息
     * exchange: topicExchange
     * bindingKey: topic.man 和 topic.# 可以接收到消息
     *
     * 创建两个 JAVA 的接收器,一个接收 topic.man,一个接收 topic.woman
     * JAVA 接收器 topic.man 可以绑定,
     *          topic.#
     *          topic.man
     * JAVA 接收器 topic.woman 只能绑定
     *          topic.#
     */
    
    @Test
    public void sendTopicMessage1() {
        sendMan();
    }

    @Test
    public void sendTopicMessage2() {
        sendWoMan();
    }

    @Test
    public void sendTopicMessage() throws InterruptedException {
        while (true) {
            // sendMan();
            sendWoMan();
            sleep(100);
        }
    }




    private void sendMan() {
        String messageId = String.valueOf(UUID.randomUUID());
        String messageData = "message: M A N ";
        String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
        Map<String, Object> manMap = new HashMap<>();
        manMap.put("messageId", messageId);
        manMap.put("messageData", messageData);
        manMap.put("createTime", createTime);
        rabbitTemplate.convertAndSend("topicExchange", "topic.man", manMap);
    }


    private void sendWoMan() {
        String messageId = String.valueOf(UUID.randomUUID());
        String messageData = "message: woman is all ";
        String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
        Map<String, Object> womanMap = new HashMap<>();
        womanMap.put("messageId", messageId);
        womanMap.put("messageData", messageData);
        womanMap.put("createTime", createTime);
        rabbitTemplate.convertAndSend("topicExchange", "topic.woman", womanMap);
    }


}

配置类 (TopicRabbitConfig.java)
@Configuration
public class TopicRabbitConfig {
    
    @Bean
    public Queue firstQueue() {
        return new Queue(RabbitMqConstant.TOPIC_ROUTING_KEY_MAN); // topic.man
    }

    @Bean
    public Queue secondQueue() {
        return new Queue(RabbitMqConstant.TOPIC_ROUTING_KEY_WOMAN); // topic.woman
    }

    @Bean
    TopicExchange exchange() {
        return new TopicExchange(RabbitMqConstant.TOPIC_EXCHANGE);
    }

    // 绑定:只接收 topic.man
    @Bean
    Binding bindingExchangeMessageByMan() {
        return BindingBuilder.bind(firstQueue()).to(exchange()).with(RabbitMqConstant.TOPIC_ROUTING_KEY_MAN);
    }

    // 绑定:接收 topic.# (所有以 topic. 开头的)
    @Bean
    Binding bindingExchangeMessageAll() {
        return BindingBuilder.bind(secondQueue()).to(exchange()).with(RabbitMqConstant.TOPIC_ROUTING_KEY_ALL);
    }
}
消费者 (TopicReceiver.java)
@Component
@RabbitListener(queues = "topic.man")
public class TopicManReceiver {

    @RabbitHandler
    public void process(Map testMessage) {
        System.out.println("【man】TopicManReceiver消费者收到消息  : " + testMessage.toString());
    }
}

@Component
@RabbitListener(queues = "topic.woman")
public class TopicTotalReceiver {

    @RabbitHandler
    public void process(Map testMessage) {
        System.out.println("【#】TopicTotalReceiver消费者收到消息  : " + testMessage.toString());
    }
}

3.3 扇形交换机 (Fanout Exchange)

  • 机制: 广播。忽略 Routing Key,发送给所有绑定的队列。
  • 场景: 系统广播、缓存失效通知。
  • 特点: 速度最快,但灵活性最低。
测试类 (FanoutSender.java)
@SpringBootTest
class FanoutTests {

    @Autowired
    RabbitTemplate rabbitTemplate;  // 使用RabbitTemplate,这提供了接收/发送等等方法
 
    @Test
    public void sendFanoutMessage() throws InterruptedException {
        while (true) {
            String messageId = String.valueOf(UUID.randomUUID());
            String messageData = "message: testFanoutMessage ";
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String, Object> map = new HashMap<>();
            map.put("messageId", messageId);
            map.put("messageData", messageData);
            map.put("createTime", createTime);
            rabbitTemplate.convertAndSend("fanoutExchange", null, map);
            sleep(100);
        }
    }

}

配置类 (FanoutRabbitConfig.java)
@Configuration
public class FanoutRabbitConfig {
    
    @Bean public Queue queueA() { return new Queue("fanout.A"); }
    @Bean public Queue queueB() { return new Queue("fanout.B"); }
    @Bean public Queue queueC() { return new Queue("fanout.C"); }

    @Bean
    FanoutExchange fanoutExchange() {
        return new FanoutExchange("fanoutExchange");
    }

    @Bean Binding bindingExchangeA() { return BindingBuilder.bind(queueA()).to(fanoutExchange()); }
    @Bean Binding bindingExchangeB() { return BindingBuilder.bind(queueB()).to(fanoutExchange()); }
    @Bean Binding bindingExchangeC() { return BindingBuilder.bind(queueC()).to(fanoutExchange()); }
}
消费者 (FanoutReceiver.java)
@Component
@RabbitListener(queues = "fanout.A")
public class FanoutReceiverA {

    @RabbitHandler
    public void process(Map testMessage) {
        System.out.println("FanoutReceiverA消费者收到消息  : " +testMessage.toString());
    }

}
@Component
@RabbitListener(queues = "fanout.B")
public class FanoutReceiverB {

    @RabbitHandler
    public void process(Map testMessage) {
        System.out.println("FanoutReceiverB消费者收到消息  : " +testMessage.toString());
    }

}
@Component
@RabbitListener(queues = "fanout.C")
public class FanoutReceiverC {

    @RabbitHandler
    public void process(Map testMessage) {
        System.out.println("FanoutReceiverC消费者收到消息  : " + testMessage.toString());
    }

}

第四部分:高级特性 - 死信与延迟队列

4.1 死信队列 (Dead Letter Exchange, DLX)

原理: 当消息在普通队列中变成“死信”时,自动转发到指定的死信交换机。
触发条件:

  1. 消息被拒绝 (basic.reject / basic.nack) 且 requeue=false
  2. 消息在队列中过期 (TTL)。
  3. 队列达到最大长度,消息被挤出。
测试类 (DeadLetterSender.java)
@SpringBootTest
public class DelayTests {
    @Autowired
    RabbitTemplate rabbitTemplate;  // 使用RabbitTemplate,这提供了接收/发送等等方法
    @Test
    public  void testDelay() throws InterruptedException { 
        while (true) {
            String string = UUID.randomUUID().toString();
            String createTime = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
            Map<String, Object> map = new HashMap<>();
            map.put("messageId",string);
            map.put("messageData","订单信息:id=1,time=2022年03月30日11:41:47");
            map.put("createTime", createTime);
            CorrelationData correlationData = new CorrelationData();

            rabbitTemplate.convertAndSend(DelayRabbitConfig.EXCHANGE_NAME,
                    "test.order.msg", map.toString().getBytes(StandardCharsets.UTF_8),correlationData);
            Thread.sleep(100);
        } 
    }
}
配置类 (DeadLetterRabbitConfig.java)
@Configuration
public class DeadLetterRabbitConfig {
    public static final String EXCHANGE_NAME_DLX = "exchange_dlx6";
    public static final String QUEUE_NAME_DLX = "queue_dlx6";
    public static final String EXCHANGE_NAME = "test_exchange_dlx6";
    public static final String QUEUE_NAME = "test_queue_dlx6";

    // 1. 普通交换机
    @Bean("test_exchange_dlx")
    public Exchange bootExchange(){
        return ExchangeBuilder.topicExchange(EXCHANGE_NAME).durable(true).build();
    }

    // 2. 死信交换机
    @Bean("exchange_dlx")
    public Exchange dlxExchange(){
        return ExchangeBuilder.topicExchange(EXCHANGE_NAME_DLX).durable(true).build();
    }

    // 3. 普通队列 (关键:设置死信参数)
    @Bean("test_queue_dlx")
    public Queue bootQueue(){
        Map<String, Object> args = new HashMap<>();
        args.put("x-dead-letter-exchange", EXCHANGE_NAME_DLX); // 死信交换机
        args.put("x-dead-letter-routing-key", "dlx.#");        // 死信路由键
        args.put("x-message-ttl", 10000);                      // 消息 TTL 10 秒
        args.put("x-max-length", 10);                          // 队列最大长度 10
        return QueueBuilder.durable(QUEUE_NAME).withArguments(args).build();
    }

    // 4. 死信队列
    @Bean("queue_dlx")
    public Queue dlxQueue(){
        return QueueBuilder.durable(QUEUE_NAME_DLX).build();
    }

    // 5. 绑定关系
    @Bean
    public Binding bindDlxQueue(@Qualifier("queue_dlx") Queue queue,
                                @Qualifier("exchange_dlx") Exchange exchange){
        return BindingBuilder.bind(queue).to(exchange).with("dlx.#").noargs();
    }
    
    @Bean
    public Binding bindNormalQueue(@Qualifier("test_queue_dlx") Queue queue,
                                   @Qualifier("test_exchange_dlx") Exchange exchange){
        return BindingBuilder.bind(queue).to(exchange).with("test.dlx.#").noargs();
    }
}

4.2 延迟队列 (Delay Queue) - 双方案详解

方案一:TTL + 死信队列 (通用方案,无需插件)

原理: 利用消息过期成为死信的机制。
缺点: 如果队列头部消息未过期,会阻塞后面已过期的消息(FIFO 限制)。
配置: 参考 DelayRabbitConfig.java (逻辑同死信配置,重点设置 x-message-ttl)。
消费者: 必须监听死信队列

@Component
public class DelayReceiver implements ChannelAwareMessageListener {
    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        try {
            String body = new String(message.getBody(), StandardCharsets.UTF_8);
            System.out.println("⏰ [" + LocalDateTime.now() + "] 接收到延迟消息: " + body);
            
            // 业务逻辑:检查订单状态,若未支付则取消
            // ...
            
            channel.basicAck(deliveryTag, false);
        } catch (Exception e) {
            // 拒绝且不重回队列,防止死循环
            channel.basicNack(deliveryTag, false, false); 
            e.printStackTrace();
        }
    }
}
方案二:RabbitMQ 延迟插件 (推荐,高性能)

原理: 安装 rabbitmq_delayed_message_exchange 插件,交换机类型变为 x-delayed-message
优点: 消息在交换机层面延迟,不阻塞队列,精度更高。
步骤:

  1. 下载对应版本的 .ez 插件放入 plugins 目录。
  2. 启用插件:rabbitmq-plugins enable rabbitmq_delayed_message_exchange
  3. 代码修改:
@Bean
public CustomExchange delayExchange() {
    Map<String, Object> args = new HashMap<>();
    args.put("x-delayed-type", "direct"); // 内部使用 direct 交换
    return new CustomExchange("delayed-exchange", "x-delayed-message", true, false, args);
}

// 发送时指定延迟时间
Map<String, Object> headers = new HashMap<>();
headers.put("x-delay", 5000); // 延迟 5000ms
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder().headers(headers).build();
rabbitTemplate.send("delayed-exchange", "routing-key", new Message("msg".getBytes(), properties));

第五部分:消息可靠性与手动确认

5.1 为什么需要手动确认?

默认 AUTO 模式下,消费者一旦收到消息,RabbitMQ 就认为已处理并删除消息。若业务逻辑报错,消息将永久丢失
Manual 模式: 只有代码显式调用 basicAck,消息才会被删除。

5.2 手动确认配置 (MessageListenerConfig.java)

@Configuration
public class MessageListenerConfig {
    @Autowired
    private CachingConnectionFactory connectionFactory;
    @Autowired
    private MyAckReceiver myAckReceiver;

    @Bean
    public SimpleMessageListenerContainer container() {
        SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory);
        container.setConcurrentConsumers(1);
        container.setMaxConcurrentConsumers(5); // 动态扩容
        container.setAcknowledgeMode(AcknowledgeMode.MANUAL); // 核心:手动确认
        container.setQueueNames("TestDirectQueue", "fanout.A", "order_que_dlx_ttl");
        container.setMessageListener(myAckReceiver);
        return container;
    }
}

5.3 统一接收者 (MyAckReceiver.java)

@Component
public class MyAckReceiver implements ChannelAwareMessageListener {
    @Override
    public void onMessage(Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        String queueName = message.getMessageProperties().getConsumerQueue();
        
        try {
            String body = new String(message.getBody(), StandardCharsets.UTF_8);
            System.out.println("📥 收到消息 [队列:" + queueName + "]: " + body);

            // === 业务逻辑 ===
            if (body.contains("error")) {
                throw new RuntimeException("模拟业务异常");
            }

            // === 成功:手动 ACK ===
            channel.basicAck(deliveryTag, false);
            System.out.println("✅ 消息确认成功");

        } catch (Exception e) {
            e.printStackTrace();
            // === 失败策略 ===
            // 方案 B: 拒绝并不再入队 (requeue=false) -> 进入死信队列 (推荐)
            channel.basicNack(deliveryTag, false, false);
            System.out.println("❌ 消息处理失败,已送入死信队列");
        }
    }
}

第六部分:控制器测试与交互流程

6.1 测试控制器 (RabbitMQController.java)

@RestController
@RequestMapping("/mq")
public class RabbitMQController {
    @Autowired
    private RabbitTemplate rabbitTemplate;

    // 1. 发送直连消息
    @GetMapping("/sendDirect")
    public String sendDirect() {
        Map<String, Object> msg = buildMsg("Hello Direct");
        rabbitTemplate.convertAndSend(RabbitMqConstant.TEST_DIRECT_EXCHANGE, 
                                      RabbitMqConstant.TEST_DIRECT_ROUTING, msg);
        return "Direct 消息已发送";
    }

    // 2. 发送主题消息
    @GetMapping("/sendTopic/{type}")
    public String sendTopic(@PathVariable String type) {
        Map<String, Object> msg = buildMsg("Hello Topic: " + type);
        String routingKey = "topic." + type;
        rabbitTemplate.convertAndSend(RabbitMqConstant.TOPIC_EXCHANGE, routingKey, msg);
        return "Topic 消息已发送: " + routingKey;
    }

    // 3. 发送广播消息
    @GetMapping("/sendFanout")
    public String sendFanout() {
        Map<String, Object> msg = buildMsg("Hello Fanout Broadcast");
        rabbitTemplate.convertAndSend("fanoutExchange", "", msg);
        return "Fanout 广播已发送";
    }

    // 4. 发送延迟消息 (触发死信)
    @GetMapping("/sendDelay")
    public String sendDelay() {
        Map<String, Object> msg = buildMsg("Order ID: 1001");
        rabbitTemplate.convertAndSend("test_exchange_dlx_ttl", "test.order.delay", msg);
        return "延迟消息已发送,预计 10 秒后处理";
    }

    private Map<String, Object> buildMsg(String data) {
        Map<String, Object> map = new HashMap<>();
        map.put("messageId", UUID.randomUUID().toString());
        map.put("messageData", data);
        map.put("createTime", LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")));
        return map;
    }
}

6.2 延迟队列交互流程图

消费者 死信队列 死信交换机 普通队列 (TTL=10s) 生产者 消费者 死信队列 死信交换机 普通队列 (TTL=10s) 生产者 消息等待 10 秒... 发送消息 (正常) 消息过期,转为死信 路由到死信队列 推送消息 手动 ACK (业务处理完成)

第七部分:生产环境最佳实践与避坑

7.1 消息可靠性保障 (三板斧)

  1. 生产者: 开启 ConfirmCallback + ReturnCallback,失败时记录 DB 并重试。
  2. Broker: 队列、交换机、消息均设置为持久化 (durable=true, deliveryMode=2)。
  3. 消费者: 使用 Manual Ack,业务成功后再 ACK。

7.2 消息重复消费 (幂等性)

网络抖动可能导致 ACK 丢失,RabbitMQ 会重发消息。

  • 解决方案: 业务层必须实现幂等性
  • 手段:
    • 数据库唯一索引: 插入前检查 message_id 是否存在。
    • Redis 原子操作: setnx key value,处理前先占位。
    • 状态机判断: 订单状态只能从 PAYING -> PAID,不能重复变。

7.3 消息积压处理

  • 临时方案: 编写临时消费者,只负责把消息取出来转发到新的 Topic,然后启动大量新消费者并行处理。
  • 根本方案: 优化消费者逻辑,增加消费者实例数量 (concurrentConsumers)。

7.4 常见异常与解决

  • MessageConversionException: 通常是生产者发了 JSON,消费者没配 Jackson2JsonMessageConverter
  • ChannelClosedException: 消费者代码抛出异常未捕获,导致通道关闭。务必在 onMessagetry-catch
  • NOT_FOUND: 队列不存在。检查 @Bean 是否被扫描。
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐