跳至正文
来两杯美式
返回

消息队列(下):Spring Boot 集成 RabbitMQ 实战

By 来两杯美式
发布于更新于

上一篇讲完安装与命令,本篇进入最贴业务的部分——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-modeACK 模式生产用手动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"

适用场景:临时响应队列——RPC 模式下,消费者把处理结果发到一个临时队列,producer 用完就关连接,队列自动消失。不要把业务队列设成 exclusive——连接一断,所有消息都没了。

六、可靠性保障三件套

生产环境的 RabbitMQ 集成,可靠性靠三件事支撑:

保障层机制失败时的兜底
生产端publisher-confirms + publisher-returnsConfirm 失败 → 重投;Return 失败 → 落库补偿
broker 端队列持久化(durable=true)+ 消息持久化(deliveryMode=2broker 重启后消息仍在
消费端手动 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-confirmspublisher-confirm-type),但所有 API(ConfirmCallback / @RabbitListener / Channel.basicAck)八年没动过。当年的代码几乎可以原样搬到 2026 的项目里。


分享这篇文章:
通过邮件分享这篇文章✓ 链接已复制
查看系列全部文章
  1. 01.消息队列开篇:为什么用 MQ + 主流选型对比
  2. 02.消息队列(中):RabbitMQ 安装与命令速查
  3. 03.消息队列(下):Spring Boot 集成 RabbitMQ 实战
  4. 04.消息队列(末):RabbitMQ 集群与高可用四种模式

上一篇
消息队列(末):RabbitMQ 集群与高可用四种模式
下一篇
消息队列(中):RabbitMQ 安装与命令速查