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);
}核心理解是:
生产者不处理业务结果,只负责发消息
队列负责暂存消息
消费者异步处理消息关键流程图
简单队列模式对于java代码:
工作队列模式 Work Queue
典型使用场景
典型场景:
- 订单异步处理
- 邮件发送
- 日志处理
- 图片压缩
- 库存扣减
- 其他耗时任务
它的核心思想是:
> 一个队列中有很多任务,多个消费者一起处理,每条消息只会被其中一个消费者消费。两个消费者都监听同一个队列
图中A处理1 3 5 B处理2 4 6这样的情况发生是因为
本质上是:
多个消费者都可用,并且处理速度接近时,RabbitMQ 轮询投递消息。
注意:
两个消费者不是各自消费一份完整消息,而是竞争消费同一个队列。
所以一条消息只会被一个消费者消费,不会被两个消费者重复消费。
prefetch与acknowledge-mode
prefetch
prefetch 表示:
每个消费者最多可以持有多少条“已经投递但还没有 ACK”的消息。
例如有两个消费者。
prefetch = 1
含义:
消费者A最多拿1条未确认消息 消费者B最多拿1条未确认消息
消费者处理完并 ACK 后,RabbitMQ 才会继续给它分配下一条消息。
注意:
prefetch 不是“每次按几条一组分发”,而是限制消费者最多持有几条未 ACK 消息。
公平分发
prefetch: 1 更容易体现公平分发。
公平分发的核心是:
谁处理完并 ACK,谁才有资格继续拿下一条消息。
如果消费者 A 很快,消费者 B 很慢:
这种情况下,快消费者会处理更多消息。
这就是 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 小于等于当前值的所有未确认消息。