跳到正文
Joeplover
学习笔记·2026-05-25·约 3 分钟阅读

RabbitMQ + Spring Boot 消息队列

Spring Boot 整合 RabbitMQ 完整示例,含声明、发送、消费和确认。

消息队列架构

RabbitMQ + Spring Boot 消息队列实战

RabbitMQ 核心模型

RabbitMQ 基于 AMQP 协议,消息流路径:Producer → Exchange → Queue → Consumer。

与 Kafka 不同,RabbitMQ 的 Exchange 概念让消息路由非常灵活。Producer 不直接发消息到队列,而是发到 Exchange,Exchange 根据绑定规则决定消息投递到哪个队列。

四种交换机

1. Direct Exchange(精确匹配)

路由键精确匹配。比如订单服务发送 order.created 消息,只有绑定了 order.created 路由键的队列能收到:

@Bean
public Binding binding(Queue queue, DirectExchange exchange) {
    return BindingBuilder.bind(queue).to(exchange).with("order.created");
}

2. Topic Exchange(模式匹配)

路由键支持通配符:# 匹配零或多个单词,* 匹配一个单词。适合多级分类:

// 消费者监听所有 order 相关消息
@RabbitListener(bindings = @QueueBinding(
    exchange = @Exchange(name = "order.topic", type = "topic"),
    value = @Queue,
    key = "order.#"
))
public void handleOrder(OrderMessage msg) { ... }

3. Fanout Exchange(广播)

消息发送到所有绑定的队列,忽略路由键。适合全局通知、缓存刷新:

@Bean
public FanoutExchange cacheRefresh() {
    return new FanoutExchange("cache.refresh");
}

4. Headers Exchange(Header 匹配)

不依赖路由键,根据消息 Header 属性匹配。灵活性最高但性能较差,实际使用较少。

Spring Boot 整合

Spring Boot 通过 spring-boot-starter-amqp 自动配置:

spring:
  rabbitmq:
    host: 192.168.8.133
    port: 5672
    username: admin
    password: admin123
    virtual-host: /
    listener:
      simple:
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 2000

我的 VM 上使用 Docker 部署 RabbitMQ 3:

docker run -d --name rabbitmq   -p 5672:5672 -p 15672:15672   -e RABBITMQ_DEFAULT_USER=admin   -e RABBITMQ_DEFAULT_PASS=admin123   rabbitmq:3-management

管理控制台通过 http://192.168.8.133:15672 访问,可以查看队列的堆积情况、消费者状态、消息速率等。

实际应用场景

在 CheckByAI 项目中,RabbitMQ 作为任务分发的核心组件:

// 生产者:提交审核任务
rabbitTemplate.convertAndSend("check.exchange", "task.submit", taskMessage);

// 消费者:处理审核任务 + 手动 ACK
@RabbitListener(queues = "check.task.queue")
public void handleTask(TaskMessage task, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
    try {
        processTask(task);
        channel.basicAck(tag, false);  // 处理成功才 ACK
    } catch (Exception e) {
        channel.basicNack(tag, false, true);  // 重试
    }
}

核心价值:

  1. 异步解耦:任务提交方不需要等待处理完成
  2. 削峰填谷:瞬时流量高峰消息积在队列中慢慢消费
  3. 死信队列:重试多次失败的消息进入 DLQ 人工处理

生产踩坑

消息堆积告警:某个消费者宕机后消息在队列中持续堆积,因为没设置队列长度上限和告警阈值。解决方案:配置队列的 x-max-length 和监控告警。

手动 ACK vs 自动 ACK:自动 ACK 如果消费者处理中崩溃,消息就丢了。生产环境必须使用手动 ACK + 重试机制。

幂等消费:RabbitMQ 的 at-least-once 语义意味着同一条消息可能被消费多次。消费端需要利用业务唯一 ID 做幂等判断,防止重复处理。