Spring Boot 整合 RabbitMQ 实战指南:从环境搭建到消息可靠性,附完整代码示例
一、环境准备
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)
✅ 死信队列与延迟队列的应用
✅ 生产环境最佳实践
掌握了这些内容,你已经能够应对绝大多数实际业务场景中的消息队列需求。建议读者动手实践,把代码跑起来,加深理解。
如果这篇文章对你有帮助,请点赞、收藏、关注支持!有问题欢迎在评论区交流讨论。
更多推荐




所有评论(0)