上一篇讲完安装与命令,本篇进入最贴业务的部分——Spring Boot 怎么集成 RabbitMQ。配置项怎么填、生产者的消息怎么确认到达、消费者怎么声明队列与手动 ACK,全部贴当年的实战代码。
一、组件封装的基本要求
不管用什么 MQ,组件封装都要满足:
- 速度快、延迟足够低
- 可靠性有保障(不丢消息、能重试)
- 异步化 + 序列化支持
- 连接池化提高性能
- 完善的补偿机制(失败重投、死信处理)
这五条是任何 MQ 客户端封装的「及格线」。Spring AMQP 帮你做了前四条,第五条要业务自己设计——比如本地消息表 + 定时补偿。
二、application.yml 配置
spring:
# Jackson 序列化配置
jackson:
date-format: yyyy-MM-dd HH:mm:ss # 时间格式
time-zone: GMT+8 # 时区
default-property-inclusion: non_null # 序列化时忽略 null 字段
rabbitmq:
host: 192.168.3.188
port: 5672
username: qimeng
password: qimeng
listener:
simple:
# 每个 Channel 持有的未确认消息总数
# 当未确认数大于该值时,消费者不会再从 broker 拉新消息(背压机制)
prefetch: 2
# 消息确认方式:auto 自动 / manual 手动 / none 无需确认
acknowledge-mode: manual
# 每个 Listener 启动时生成的最小 Channel 数(并发消费者数量)
concurrency: 5
max-concurrency: 5
# 发送方确认模式
publisher-confirms: true # 消息到达 broker 时回调 ConfirmCallback
publisher-returns: true # 消息找不到队列时回调 ReturnCallback
template:
mandatory: true # 找不到队列时返回给 producer,而不是直接丢弃
2.1 关键参数解释
| 参数 | 含义 | 调优建议 |
|---|---|---|
prefetch | 每个 Channel 的预取未确认消息数 | 消费慢就调小(如 1),消费快可调大(5~20) |
acknowledge-mode | ACK 模式 | 生产用手动(manual),保证消费失败能重投 |
concurrency | 最小并发消费者数 | 按下游能力调,不是越大越好 |
max-concurrency | 最大并发消费者数 | 弹性扩容上限 |
publisher-confirms | 发送确认回调 | 生产必开,确认消息到达 broker |
publisher-returns | 路由失败回调 | 生产必开,捕获「消息找不到队列」 |
新手最大的坑:默认
acknowledge-mode: auto,业务方法抛异常会无限重试——队列瞬间堆爆。生产环境必须改 manual,把 ACK 时机掌握在自己手里。
2.2 prefetch 与背压
prefetch 是消费者侧的背压机制:
prefetch = 2
┌── 未确认 msg1 ──┐
Channel ┤── 未确认 msg2 ──┤ → broker 不会再推新消息
└────────────────┘
等到 msg1 ACK 后,broker 才推下一条
prefetch设得过大,消费者内存吃紧、消息堆积在自己进程里;设得过小,吞吐量上不去。典型起步值 2~5,按业务慢慢调。
三、生产者:发送确认(ConfirmCallback)
publisher-confirms: true 开启后,每条消息到达 broker 都会触发 ConfirmCallback:
@Component
public class RabbitProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
// Confirm 回调:消息是否成功到达 broker
final RabbitTemplate.ConfirmCallback confirmCallback = (correlationData, ack, cause) -> {
String userId = correlationData.getId();
if (ack) {
// 消息成功发送到 broker 并妥善保存
System.out.println(userId + " --> ack :" + ack);
} else {
// 发送失败:记录日志、重投、告警
System.out.println(userId + " --> nack, cause: " + cause);
}
};
@PostConstruct
public void init() {
rabbitTemplate.setConfirmCallback(confirmCallback);
}
@RequestMapping("addUser")
public String addMsg(@RequestBody User user) {
CorrelationData correlationData = new CorrelationData(user.getId());
// 参数:交换机、routing key、消息体、关联 ID
rabbitTemplate.convertAndSend("test.exchange", "aaa.bbb.ccc", user, correlationData);
return "success";
}
}
3.1 ConfirmCallback vs ReturnCallback
两个回调触发条件不同,新人很容易搞混:
| 回调 | 触发条件 | 失败语义 |
|---|---|---|
ConfirmCallback (ack=true) | 消息到达 Exchange | ✅ |
ConfirmCallback (ack=false) | Exchange 接收失败(如交换机不存在) | ❌ 在 producer 侧丢 |
ReturnCallback | 消息到 Exchange 但找不到 Queue(路由失败) | ❌ 在 broker 侧丢 |
Producer → Exchange(ConfirmCallback 触发)
│
├─ 找到 Queue → 投递成功
└─ 找不到 Queue → ReturnCallback 触发(mandatory=true 时)
生产环境两个回调都要监听——光监听 Confirm 不够,路由失败的 Return 才是消息丢失的高发地。
四、消费者:@RabbitListener 声明式订阅
最优雅的写法是注解里直接声明队列、交换机、绑定——客户端首次启动时自动创建:
@Component
public class RabbitConsumer {
/**
* 1. 每当新的消费者启动,会建立若干新 Channel(数量由 concurrency 决定)
* 每个 Channel 维护一个从 1 开始递增的 deliveryTag
*
* 2. 队列声明 exclusive = "true" 表示排他队列:
* 只对首次声明它的 Connection 可见,Connection 关闭时自动删除
*
* 3. routing key 模式匹配:
* - aaa.* 只能匹配 aaa.xxx(一段)
* - aaa.# 能匹配 aaa.xxx.xxx(多段,含 aaa 本身)
*
* 4. 队列和交换机可以有多重 binding(不同 routing key 绑同一个队列)
*/
@RabbitListener(bindings = @QueueBinding(
value = @Queue(value = "${mq.queue.test}", durable = "true"),
exchange = @Exchange(name = "${mq.exchange.test}", durable = "true", type = "topic"),
key = "aaa.#"
))
public void onMessage(@Payload User user,
Channel channel,
Message message) throws Exception {
System.out.println(user);
long deliveryTag = message.getMessageProperties().getDeliveryTag();
// 手动 ACK(第二个参数 multiple=false 表示只确认当前这条)
channel.basicAck(deliveryTag, false);
}
}
4.1 路由键通配符
| 模式 | 匹配 |
|---|---|
aaa.bbb.ccc | 精确匹配(Direct 风格) |
aaa.* | 匹配 aaa.任意一段(如 aaa.bbb,不匹配 aaa.bbb.ccc) |
aaa.# | 匹配 aaa.任意多段(如 aaa.bbb.ccc.ddd,也匹配 aaa 本身) |
*和#的差别是新人最容易记错的:*是一段,#是零到多段。
4.2 手动 ACK 的三种姿势
// 1. 业务成功 → ACK(消息从队列删除)
channel.basicAck(deliveryTag, false);
// 2. 业务异常、可重试 → NACK + requeue(消息回到队列,重新投递)
channel.basicNack(deliveryTag, false, true);
// 3. 业务异常、不可恢复 → reject(直接丢弃,或路由到死信队列)
channel.basicReject(deliveryTag, false);
| 操作 | 消息去向 | 适用场景 |
|---|---|---|
basicAck | 删除 | 正常处理完成 |
basicNack(requeue=true) | 回队列头部,重新投递 | 临时异常(如下游服务抖动) |
basicNack(requeue=false) | 进死信队列(如配置了 DLX) | 永久异常(如格式错误) |
生产建议:NACK + requeue 要配重试次数上限,否则消息会被无限重投,最终把队列打爆。可以用
RetryTemplate或自己计数。
五、排他队列(Exclusive Queue)
声明队列时设 exclusive = "true":
- 该队列只对首次声明它的 Connection 可见
- Connection 关闭时自动删除
适用场景:临时响应队列——RPC 模式下,消费者把处理结果发到一个临时队列,producer 用完就关连接,队列自动消失。不要把业务队列设成 exclusive——连接一断,所有消息都没了。
六、可靠性保障三件套
生产环境的 RabbitMQ 集成,可靠性靠三件事支撑:
| 保障层 | 机制 | 失败时的兜底 |
|---|---|---|
| 生产端 | publisher-confirms + publisher-returns | Confirm 失败 → 重投;Return 失败 → 落库补偿 |
| broker 端 | 队列持久化(durable=true)+ 消息持久化(deliveryMode=2) | broker 重启后消息仍在 |
| 消费端 | 手动 ACK + 死信队列 | 业务失败 → NACK → 死信 → 人工介入 |
三层都做到了,才能叫「消息不丢」。少一层,丢消息的概率从「理论零」变成「实际有」。
七、小结
| 角色 | 关键配置 / API |
|---|---|
| 配置 | prefetch / acknowledge-mode=manual / publisher-confirms=true |
| 生产者 | ConfirmCallback + ReturnCallback + CorrelationData |
| 消费者 | @RabbitListener + @QueueBinding + 手动 basicAck/basicNack |
| 可靠性 | 持久化 + 手动 ACK + 死信队列 |
下一篇「消息队列(末):RabbitMQ 集群与高可用」收尾——单点 RabbitMQ 在生产就是定时炸弹,必须上集群。讲清四种集群模式:普通集群、镜像队列、Warren(主备)、Shovel(远程)。
现代视角补一句(2026):到 Spring Boot 2.x/3.x,参数名从 kebab-case 改成了 camelCase(
publisher-confirms→publisher-confirm-type),但所有 API(ConfirmCallback/@RabbitListener/Channel.basicAck)八年没动过。当年的代码几乎可以原样搬到 2026 的项目里。