消息队列
MQ的本质
MQ三大核心场景
- 解耦 在传统的同步调用中,系统A若要传递数据给系统B和C,必须在A的代码中显式调用B和C的接口。如果未来系统D也需要该数据,A 的代码必须修改并重新发布。引入MQ后A变成了生产者,只负责将数据丢入MQ;B、C、D 作为消费者去MQ订阅数据。A与下游系统在物理和逻辑上彻底解绑。
- 异步
核心写库完成后直接将事件发给MQ,其他操作则在后台异步拉取消息进行处理,将原来的串行事件优化为并行事件,提高主链路的吞吐量和响应速度。 - 削峰 对于巨量请求,请求都先暂存在MQ中,后端的订单服务根据自身的极限处理能力,按照固定节奏从MQ中拉取并处理请求。
引入MQ的代价
系统可用性降低:一旦MQ集群宕机,上下游的通信就会中断,导致整个业务不可用。
系统复杂度增加:对于网络抖动导致的信息丢失、消费者重复消费、以及如何保证某些特定消息的顺序执行。
系统维护成本增加:监控队列深度、磁盘水位和内存使用率。
基本通信模型
点对点模型:生产者发送一条消息到队列,多个消费者可以监听同一个队列,但一条消息只能被一个消费者成功拉取并处理,消费完毕后,消息从队列中销毁。
发布/订阅模型:生产者将消息发送到一个主题或交换机。多个消费者独立订阅该主题。一条消息会被复制或广播,所有订阅了该主题的消费者都能接收到这条消息的完整副本。
基本例子
引入spring-boot-starter-amqp依赖项,配置队列、交换机、路由KEY,均注册为Bean。
package com.example.Config;
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.amqp.core.Queue;
@Configuration
public class RabbitConfig {
public static final String QUEUE_NAME = "demo.queue";
public static final String EXCHANGE_NAME = "demo.direct.exchange";
public static final String ROUTING_KEY = "demo.routing.key";
@Bean
public Queue demoQueue(){
return new Queue(QUEUE_NAME, true);
}
@Bean
public DirectExchange demoExchange(){
return new DirectExchange(EXCHANGE_NAME);
}
@Bean
public Binding bindingDemo(Queue demoQueue,DirectExchange demoExchange){
return BindingBuilder.bind(demoQueue).to(demoExchange).with(ROUTING_KEY);
}
}生产者注册
package com.example.Producer;
import com.example.Config.RabbitConfig;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
@Component
public class MessageProducer {
private final RabbitTemplate rabbitTemplate;
public MessageProducer(RabbitTemplate rabbitTemplate){
this.rabbitTemplate = rabbitTemplate;
}
public void sendTestMessage(String message){
rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_NAME,RabbitConfig.ROUTING_KEY,message);
System.out.println("生产者,消息已投递到交换机");
}
}消费者监听队列
package com.example.Consumer;
import com.example.Config.RabbitConfig;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
public class MessageConsumer {
@RabbitListener(queues = RabbitConfig.QUEUE_NAME)
public void receiveMessage(String message) {
System.out.println("消费者从队列中拉取到信息" + message);
}
}进阶
消息不丢失
- 发送端丢失风险(生产者->Broker) 由于网络抖动,消息在半路丢失;交换机/路由键配置错误,导致Broker收到消息但无法将其路由到任何队列。
解决办法:引入确认机制与重试机制。让生产者能够监听Broker的回执,明确知道消息是否成功到达交换机,以及是否成功到路由到路由。如果这一步失败,需要触发重试逻辑。 - Broker端丢失风险(MQ内部存储)
消息成功到达队列,暂存在内存中,如果RabbitMQ节点突然断电或进程崩溃,内存中的未落盘数据全部丢失。
解决办法:配置消息持久化机制。这要求,要求交换机、队列、消息实体三者都必须标记为持久化。 - 消费端丢失风险(Broker->消费者)
消费者从队列中拉取消息,SpringBoot默认采用自动确认(Auto ACK),Broker见其被拉取便将其从队列中删除。但消费者在执行业务逻辑时发生异常或宕机。业务未完成,且消息已经物理销毁了。
解决办法:开启手动ACK操作。关闭自动确认,要求消费者在业务代码彻底执行成功后,才通过代码显式地向Broker发送确认指令;若执行失败,则显式拒绝该消息,并根据策略将其重新入队或丢弃。
死信队列
channel.basicNack(tag,false,true)会将处理失败的消息重新放回原队列,如果这个消息处理业务始终失败,会导致循环阻塞。然而信息又最好不消失,可以让RabbitMQ自动将它转移到一个专门的收容队列,等待后续人工或定时任务处理,这个收容队列就是死信队列。
消息不重复(幂等性设计)
如果因为网络抖动或者消费者在执行完业务逻辑后,发送ACK给Broker之前突然宕机,RabbitMQ会认为这条消息没有被重复消费,从而重新投递,同一笔订单可能处理多次。实现幂等性设计,无论一条消息被执行多少次,产生的影响和第一次执行完全相同,结合Redis的分布式锁或者MySQL的唯一索引来实现。在业务处理前先查看,再执行。
消息顺序性保证
RabbitMQ本身是一个先进先出的队列,单个队列内部是保证顺序的。但当多消费者监听同一个队列,由于执行速度不同就有可能乱序;消息1处理失败重新回到队列重试,而消息2被成功处理,同样乱序。通常要求的顺序是局部顺序而不是全局顺序,通过局部顺序映射来解决,将需要保证顺序的同一批消息,通过一致性Hash路由到同一个固定的Queue中,且该Queue只允许由一个单线程的Comsumer去消费。
高速与高可用
高速
消息队列采用磁盘顺序追加写,新的消息队列只在文件末尾追加,不修改旧数据。MQ同样采用页缓存,接收到的消息先写入PageCache,操作系统会在后台按时将PageCache中的数据批量刷入磁盘。MQ采用零拷贝技术,直接在内核态将数据从Read Buffer传输到Socket Buffer,完全跳过了用户空间,加快了数据传输速度。
高可用
为了应对单机宕机,生产环境部署集群。
- 普通集群:RabbitMQ默认集群,多台机器共享元数据,但消息实体只存在于创建该队列的节点上。其他节点只负责路由,如果存储消息的节点宕机,消息依然不可用。
- 镜像/多副本集群:同一条消息被复制到集群内的多台机器上。即使主节点直接损坏,副本依然包含完整数据。
从节点消息同步有两种策略:
- 同步双写:主节点必须等从节点也写入成功,才向生产者返回ACK。虽然数据绝对不丢失,但是延迟翻倍,吞吐量下降。
- 异步复制:主节点写完立刻返回ACK,后台再同步给从节点。吞吐量高,但如果ACK后,同步前宕机,会造成少量数据丢失。 当主节点宕机了,集群内的节点通过心跳机制互相探测,判定主节点失联后,触发故障转移。根据共识算法,拥有最新、最全数据的从节点通常会被选为新的主节点。