Administrator
发布于 2026-06-08 / 7 阅读
0

RabbitMq-01

RabbitMQ 基础模型

RabbitMQ 是一个消息队列中间件。它的作用是让系统之间不要直接强依赖调用,而是通过“发送消息”和“消费消息”的方式完成协作。

同步调用是:

A 系统直接调用 B 系统

如果 B 系统慢了、挂了,A 系统也会受到影响。

使用 RabbitMQ 后变成:

A 系统把消息发送到 RabbitMQ
B 系统从 RabbitMQ 中消费消息

这样 A 和 B 不需要强绑定在一起。

RabbitMQ 中几个核心角色:

Producer:生产者,负责发送消息
Consumer:消费者,负责接收并处理消息
Broker:RabbitMQ 服务本身
Exchange:交换机,负责接收生产者的消息并决定发给哪个队列
Queue:队列,真正存储消息的地方
Binding:绑定关系,连接 Exchange 和 Queue
RoutingKey:路由键,用来决定消息进入哪个队列

消息为什么不直接发给队列,而要经过交换机?

因为 RabbitMQ 设计上是:

生产者不直接关心队列
生产者只把消息发给交换机
交换机根据规则把消息路由到队列
消费者只监听队列

这样可以支持更灵活的消息分发方式,比如一个消息发给一个队列、多个队列,或者根据不同 routingKey 发给不同队列。

常见使用场景:

异步处理:下单后发送短信、邮件,不阻塞主流程
削峰填谷:高并发请求先进入队列,后端慢慢消费
系统解耦:订单系统不直接依赖库存、通知系统
失败重试:消费失败后可以重新投递或进入死信队列

简单队列模式

简单队列模式是 RabbitMQ 最小的消息模型:

一个生产者 → 一个队列 → 一个消费者

交换机名称是空字符串 ""

在这个模式里,通常使用的是 RabbitMQ 的默认交换机。

默认交换机的特点:

交换机名称是空字符串 ""
RoutingKey 必须等于队列名称
消息会被路由到同名队列

消息流是:

生产者发送消息
RabbitMQ 根据 routingKey 找到同名队列
消息进入 simple.queue
消费者监听 simple.queue
消费者收到并处理消息

Spring Boot 中代码一般对应关系是:

RabbitTemplate:负责发送消息
@RabbitListener:负责监听队列
Queue:声明队列
DTO / Request 对象:作为消息体

发送端大概是:

rabbitTemplate.convertAndSend("simple.queue", message);

这里如果只传两个参数,本质上就是使用默认交换机:

exchange = ""
routingKey = "simple.queue"

消费者大概是:

@RabbitListener(queues = "simple.queue")
public void listen(String message) {
    System.out.println("收到消息:" + message);
}

核心理解是:

生产者不处理业务结果,只负责发消息
队列负责暂存消息
消费者异步处理消息

关键流程图

flowchart LR A["Producer 生产者"] --> B["Exchange 交换机"] B --> C["Queue 队列"] C --> D["Consumer 消费者"] A -. "发送消息 + RoutingKey" .-> B B -. "根据 Binding / RoutingKey 路由" .-> C D -. "监听队列并消费" .-> C

简单队列模式对于java代码:

flowchart LR A["Spring Boot 接口<br/>发送消息"] --> B["RabbitTemplate"] B --> C["默认交换机<br/>exchange = 空字符串"] C --> D["simple.queue"] D --> E["@RabbitListener<br/>消费者监听"]

工作队列模式 Work Queue

典型使用场景

典型场景:

- 订单异步处理
- 邮件发送
- 日志处理
- 图片压缩
- 库存扣减
- 其他耗时任务
它的核心思想是:

> 一个队列中有很多任务,多个消费者一起处理,每条消息只会被其中一个消费者消费。

两个消费者都监听同一个队列

图中A处理1 3 5 B处理2 4 6这样的情况发生是因为

本质上是:

多个消费者都可用,并且处理速度接近时,RabbitMQ 轮询投递消息。

flowchart TD Q["工作队列<br/>消息1 消息2 消息3 消息4 消息5 消息6"] --> C1["消费者 A"] Q --> C2["消费者 B"] C1 --> A1["处理消息1"] C1 --> A2["处理消息3"] C1 --> A3["处理消息5"] C2 --> B1["处理消息2"] C2 --> B2["处理消息4"] C2 --> B3["处理消息6"]

注意:

两个消费者不是各自消费一份完整消息,而是竞争消费同一个队列。

所以一条消息只会被一个消费者消费,不会被两个消费者重复消费。

prefetch与acknowledge-mode

prefetch

prefetch 表示:

每个消费者最多可以持有多少条“已经投递但还没有 ACK”的消息。

例如有两个消费者。

prefetch = 1

flowchart LR Q["队列"] --> A["消费者A<br/>最多1条未ACK"] Q --> B["消费者B<br/>最多1条未ACK"]

含义:

消费者A最多拿1条未确认消息 消费者B最多拿1条未确认消息

消费者处理完并 ACK 后,RabbitMQ 才会继续给它分配下一条消息。

注意:

prefetch 不是“每次按几条一组分发”,而是限制消费者最多持有几条未 ACK 消息。

公平分发

prefetch: 1 更容易体现公平分发。

公平分发的核心是:

谁处理完并 ACK,谁才有资格继续拿下一条消息。

如果消费者 A 很快,消费者 B 很慢:

sequenceDiagram participant Q as 工作队列 participant A as 快消费者 participant B as 慢消费者 Q->>A: 消息1 Q->>B: 消息2 A-->>Q: ACK消息1 Q->>A: 消息3 A-->>Q: ACK消息3 Q->>A: 消息4 B-->>Q: ACK消息2 Q->>B: 消息5

这种情况下,快消费者会处理更多消息。

这就是 prefetch: 1 的意义:
慢消费者不会一次性拿走太多消息,快消费者可以持续处理新任务。

acknowledge-mode

当前配置是手动 ACK:

acknowledge-mode: manual

消费者处理完成后,需要手动确认:

channel.basicAck(deliveryTag, false);

第二个参数 false 表示:

只确认当前这一条消息。

所以即使配置:

prefetch: 2

也不是“消费两条后再 ACK”。

实际逻辑仍然是:

最多可以提前拿2条未确认消息 但是代码每处理完1条,就ACK当前这一条

如果写成:

channel.basicAck(deliveryTag, true);

才表示批量 ACK,会确认当前通道上 deliveryTag 小于等于当前值的所有未确认消息。