一、环境准备
1.1 使用 Docker 启动 RabbitMQ
推荐使用 Docker 快速启动 RabbitMQ 环境,一条命令搞定:

docker run -d --name rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  rabbitmq:3.12-management

启动后访问 http://localhost:15672,使用默认账号 guest/guest 登录管理界面。

1.2 创建 Spring Boot 项目
在 pom.xml 中添加 Spring AMQP 依赖:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>```
Spring AMQP 是 Spring 官方提供的 RabbitMQ 集成工具,基于 Spring Boot 实现了自动装配,使用起来非常方便。它提供了三个核心能力:自动声明队列和交换机、基于注解的监听器模式、封装了 RabbitTemplate 工具类用于发送消息。

1.3 配置文件
在 application.yml 中配置 RabbitMQ 连接信息:

```yaml
spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    virtual-host: /
    # 生产者确认
    publisher-confirm-type: correlated
    publisher-returns: true
    # 消费者确认
    listener:
      simple:
        acknowledge-mode: manual

二、五种工作模式的完整实现
RabbitMQ 支持五种经典的消息模式,下面我们逐一实现。

2.1 简单模式(Simple Queue)
最基础的模式:一个生产者、一个消费者,一对一通信。

配置类:

@Configuration
public class SimpleQueueConfig {
    @Bean
    public Queue simpleQueue() {
        return new Queue("simple.queue");
    }
}```
生产者:

```java
@RestController
public class SimpleProducer {
    @Autowired
    private RabbitTemplate rabbitTemplate;

    @GetMapping("/send/simple")
    public String sendSimple(@RequestParam String msg) {
        rabbitTemplate.convertAndSend("simple.queue", msg);
        return "消息发送成功:" + msg;
    }
}

消费者:

@Component
public class SimpleConsumer {
    @RabbitListener(queues = "simple.queue")
    public void receive(String msg) {
        System.out.println("收到消息:" + msg);
    }
}

2.2 工作队列模式(Work Queues)
一个生产者、多个消费者,消息在消费者之间竞争消费,适用于任务分发场景。

@Configuration
public class WorkQueueConfig {
    @Bean
    public Queue workQueue() {
        return new Queue("work.queue");
    }
}


@Component
public class WorkConsumer {
    // 两个消费者竞争消费同一队列的消息
    @RabbitListener(queues = "work.queue")
    public void receive1(String msg) throws InterruptedException {
        System.out.println("消费者1收到:" + msg);
        Thread.sleep(1000); // 模拟处理耗时
    }

    @RabbitListener(queues = "work.queue")
    public void receive2(String msg) throws InterruptedException {
        System.out.println("消费者2收到:" + msg);
        Thread.sleep(2000); // 模拟处理耗时
    }
}

注意:默认情况下 RabbitMQ 采用轮询分发机制。如果要实现“能者多劳”,需要设置 prefetch=1,让处理能力强的消费者获得更多消息。

2.3 发布订阅模式(Publish/Subscribe)
使用 Fanout Exchange 将消息广播到所有绑定的队列。

@Configuration
public class FanoutConfig {
    @Bean
    public FanoutExchange fanoutExchange() {
        return new FanoutExchange("fanout.exchange");
    }

    @Bean
    public Queue fanoutQueue1() {
        return new Queue("fanout.queue1");
    }

    @Bean
    public Queue fanoutQueue2() {
        return new Queue("fanout.queue2");
    }

    @Bean
    public Binding binding1() {
        return BindingBuilder.bind(fanoutQueue1()).to(fanoutExchange());
    }

    @Bean
    public Binding binding2() {
        return BindingBuilder.bind(fanoutQueue2()).to(fanoutExchange());
    }
}

2.4 路由模式(Routing)
使用 Direct Exchange,根据 Routing Key 精确匹配路由消息。

java

@Configuration
public class DirectConfig {
    @Bean
    public DirectExchange directExchange() {
        return new DirectExchange("direct.exchange");
    }

    @Bean
    public Queue errorQueue() {
        return new Queue("error.queue");
    }

    @Bean
    public Queue infoQueue() {
        return new Queue("info.queue");
    }

    @Bean
    public Binding errorBinding() {
        return BindingBuilder.bind(errorQueue())
                .to(directExchange()).with("error");
    }

    @Bean
    public Binding infoBinding() {
        return BindingBuilder.bind(infoQueue())
                .to(directExchange()).with("info");
    }
}

发送消息时指定 Routing Key:

java

rabbitTemplate.convertAndSend("direct.exchange", "error", "系统异常消息");
rabbitTemplate.convertAndSend("direct.exchange", "info", "系统通知消息");

2.5 主题模式(Topics)
使用 Topic Exchange,支持通配符匹配(* 匹配一个单词,# 匹配零个或多个单词)。

java

@Configuration
public class TopicConfig {
    @Bean
    public TopicExchange topicExchange() {
        return new TopicExchange("topic.exchange");
    }

    @Bean
    public Queue orderQueue() {
        return new Queue("order.queue");
    }

    @Bean
    public Queue allQueue() {
        return new Queue("all.queue");
    }

    @Bean
    public Binding orderBinding() {
        // 只接收订单相关消息
        return BindingBuilder.bind(orderQueue())
                .to(topicExchange()).with("order.#");
    }

    @Bean
    public Binding allBinding() {
        // 接收所有消息
        return BindingBuilder.bind(allQueue())
                .to(topicExchange()).with("#");
    }
}

三、消息可靠性保障实战
3.1 生产者确认机制(Publisher Confirm)
开启 Confirm 模式,确保消息成功到达 Broker:

java

@Component
public class ConfirmCallbackService implements RabbitTemplate.ConfirmCallback {

    @Override
    public void confirm(CorrelationData correlationData, boolean ack, String cause) {
        if (ack) {
            System.out.println("消息成功到达 Exchange,ID:" + correlationData.getId());
        } else {
            System.out.println("消息发送失败,原因:" + cause);
            // 此处可以加入重发逻辑或记录失败消息
        }
    }
}

同时在配置类中设置:

java

@PostConstruct
public void init() {
    rabbitTemplate.setConfirmCallback(confirmCallbackService);
}

3.2 消息持久化
java

@Bean
public Queue durableQueue() {
    // durable: true 表示队列持久化
    return QueueBuilder.durable("durable.queue").build();
}

// 发送持久化消息
Message message = MessageBuilder
    .withBody("持久化消息".getBytes())
    .setDeliveryMode(MessageDeliveryMode.PERSISTENT) // 消息持久化
    .build();
rabbitTemplate.send("durable.queue", message);

3.3 消费者手动确认
java

@Component
public class AckConsumer {
    @RabbitListener(queues = "durable.queue")
    public void receive(String msg, Channel channel, Message message) throws IOException {
        try {
            System.out.println("处理消息:" + msg);
            // 处理成功,手动确认
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        } catch (Exception e) {
            // 处理失败,拒绝并重新入队
            channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
        }
    }
}

四、高级特性:死信队列与延迟队列
4.1 死信队列(Dead Letter Queue)
死信队列用于处理无法被正常消费的消息,常见的“死信”来源包括:

消息被消费者拒绝(basic.reject/basic.nack)且 requeue=false

消息 TTL(生存时间)到期

队列达到最大长度

java

@Configuration
public class DeadLetterConfig {

    // 死信交换机
    @Bean
    public DirectExchange deadLetterExchange() {
        return new DirectExchange("dlx.exchange");
    }

    // 死信队列
    @Bean
    public Queue deadLetterQueue() {
        return new Queue("dlx.queue");
    }

    @Bean
    public Binding dlxBinding() {
        return BindingBuilder.bind(deadLetterQueue())
                .to(deadLetterExchange()).with("dlx.routing.key");
    }

    // 普通业务队列(绑定死信交换机)
    @Bean
    public Queue businessQueue() {
        return QueueBuilder.durable("business.queue")
                .deadLetterExchange("dlx.exchange")
                .deadLetterRoutingKey("dlx.routing.key")
                .ttl(10000) // 消息10秒后过期
                .maxLength(100) // 队列最大长度
                .build();
    }
}

4.2 延迟队列
RabbitMQ 本身不直接支持延迟队列,但可以通过 “TTL + 死信队列” 组合间接实现延迟消费的效果。消息在普通队列中到期后,自动转发到死信队列,消费者监听死信队列即可实现延迟处理。

典型应用场景:订单超时取消(30分钟未支付自动取消)、定时任务触发。

java

@Configuration
public class DelayQueueConfig {

    @Bean
    public Queue delayQueue() {
        return QueueBuilder.durable("delay.queue")
                .deadLetterExchange("dlx.exchange")
                .deadLetterRoutingKey("dlx.routing.key")
                .ttl(30000) // 30秒延迟
                .build();
    }

    // 发送延迟消息
    public void sendDelayMsg(String msg) {
        rabbitTemplate.convertAndSend("delay.queue", (Object) msg);
    }

    // 消费者监听死信队列
    @RabbitListener(queues = "dlx.queue")
    public void handleDelayMsg(String msg) {
        System.out.println("延迟消息已处理:" + msg);
    }
}

五、生产环境最佳实践
5.1 队列设计规范
队列命名:采用 {业务模块}.{功能}.queue 格式,如 order.pay.queue

交换机命名:采用 {业务模块}.{类型}.exchange 格式,如 order.topic.exchange

Routing Key 命名:采用点分格式,如 order.create.success

5.2 性能优化要点
合理使用连接池:复用 Connection 和 Channel,避免频繁创建销毁

设置合适的 Prefetch:平衡吞吐量和消费者负载

避免消息堆积:设置队列长度限制和 TTL,防止无限堆积导致内存溢出

使用惰性队列:将消息尽可能存储到磁盘,减少内存占用,解决消息堆积问题

5.3 监控与告警
关注关键指标:队列深度、消息速率、消费者数量、连接数和通道数,通过 Prometheus + Grafana 搭建监控面板,设置阈值告警。

总结
本文从环境搭建开始,逐步深入,完整覆盖了 Spring Boot 整合 RabbitMQ 的核心内容:

✅ 五种工作模式的完整代码实现

✅ 消息可靠性保障(Confirm、持久化、ACK)

✅ 死信队列与延迟队列的应用

✅ 生产环境最佳实践

掌握了这些内容,你已经能够应对绝大多数实际业务场景中的消息队列需求。建议读者动手实践,把代码跑起来,加深理解。

如果这篇文章对你有帮助,请点赞、收藏、关注支持!有问题欢迎在评论区交流讨论。

Logo

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

更多推荐