RabbitMQ消息队列:异步解耦与业务削峰
📅 2026/7/31 1:07:13
👁️ 阅读次数
📝 编程学习
RabbitMQ消息队列:异步解耦与业务削峰
同步调用就像你打电话等对方接——对方不接你就一直卡着;异步消息就像发微信——发完该干嘛干嘛,对方有空了自然回你。
一、消息队列解决了什么问题
在单体架构时代,所有功能揉在一个项目里,方法之间直接调用,简单粗暴。但一旦系统变大,问题就来了:
- 异步处理:用户注册后要发邮件、发短信、发优惠券……同步调用的话用户得等半天,体验极差。丢到消息队列里,注册接口秒回,后续操作慢慢消费。
- 应用解耦:订单系统直接调用库存系统,库存挂了订单也跟着挂。中间加个队列,订单只管发消息,库存恢复了继续消费即可。
- 流量削峰:秒杀场景瞬间涌入10万请求,数据库直接被干趴。队列做个缓冲,消费者按自己的节奏处理,系统稳如老狗。
- 日志收集:分布式系统中各服务把日志推到队列,由统一的日志服务消费存储,EFK/ELK的经典套路。
二、RabbitMQ核心概念
RabbitMQ的消息流转模型如下:
Producer → Exchange → (Binding) → Queue → Consumer 生产者 交换机 绑定 队列 消费者- Producer(生产者):产生消息的应用程序
- Exchange(交换机):接收生产者发送的消息,根据路由规则分发到队列
- Queue(队列):存放消息的缓冲区,消息在这里排队等消费
- Binding(绑定):交换机和队列之间的关联关系,附带路由键
- Consumer(消费者):从队列中获取消息并处理的应用程序
三、交换机四种类型
RabbitMQ提供了四种Exchange类型,理解清楚就知道消息怎么路由了。
3.1 Direct(直连)
最简单的模式,消息的路由键(routing key)和绑定的键完全匹配,消息才会被投递到对应队列。
routing key = "order.create" → 只匹配绑定 "order.create" 的队列3.2 Fanout(扇出)
广播模式,忽略路由键,消息被投递到与该交换机绑定的所有队列。适合广播通知场景。
3.3 Topic(主题)
支持通配符匹配,灵活性最高:
*匹配一个单词#匹配零个或多个单词
绑定键 "order.*" → 匹配 "order.create"、"order.cancel",不匹配 "order.create.detail" 绑定键 "order.#" → 匹配 "order.create"、"order.create.detail" 全都匹配3.4 Headers(头部)
不靠路由键,而是根据消息头(headers)中的键值对匹配。用的少,了解即可。
四、SpringBoot整合RabbitMQ
4.1 引入依赖
<dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-amqp</artifactId></dependency>4.2 yml配置
spring:rabbitmq:host:127.0.0.1port:5672username:guestpassword:guest# 消息确认机制publisher-confirm-type:correlated# 发布确认publisher-returns:true# 消息返回listener:simple:acknowledge-mode:manual# 手动ACKprefetch:1# 每次拉取消息数4.3 队列与交换机配置
@ConfigurationpublicclassRabbitMQConfig{// 队列名称publicstaticfinalStringEMAIL_QUEUE="email.queue";publicstaticfinalStringSMS_QUEUE="sms.queue";publicstaticfinalStringORDER_EXCHANGE="order.exchange";publicstaticfinalStringORDER_ROUTING_KEY="order.notify";@BeanpublicDirectExchangeorderExchange(){returnnewDirectExchange(ORDER_EXCHANGE,true,false);}@BeanpublicQueueemailQueue(){returnnewQueue(EMAIL_QUEUE,true);}@BeanpublicQueuesmsQueue(){returnnewQueue(SMS_QUEUE,true);}@BeanpublicBindingemailBinding(QueueemailQueue,DirectExchangeorderExchange){returnBindingBuilder.bind(emailQueue).to(orderExchange).with(ORDER_ROUTING_KEY);}@BeanpublicBindingsmsBinding(QueuesmsQueue,DirectExchangeorderExchange){returnBindingBuilder.bind(smsQueue).to(orderExchange).with(ORDER_ROUTING_KEY);}}五、发送消息:RabbitTemplate
@ServicepublicclassOrderService{@AutowiredprivateRabbitTemplaterabbitTemplate;publicvoidcreateOrder(OrderDTOorderDTO){// 1. 保存订单(数据库操作省略)// ...// 2. 异步发送通知消息Stringmsg=JSON.toJSONString(orderDTO);rabbitTemplate.convertAndSend(RabbitMQConfig.ORDER_EXCHANGE,RabbitMQConfig.ORDER_ROUTING_KEY,msg);// 3. 直接返回,不等邮件/短信发送完成return;}}六、接收消息:@RabbitListener
@ComponentpublicclassEmailConsumer{@RabbitListener(queues=RabbitMQConfig.EMAIL_QUEUE)@RabbitHandlerpublicvoidreceive(Stringmessage,Channelchannel,MessagemessageObj)throwsIOException{longdeliveryTag=messageObj.getMessageProperties().getDeliveryTag();try{OrderDTOorder=JSON.parseObject(message,OrderDTO.class);// 发送邮件逻辑System.out.println("发送邮件到:"+order.getEmail());// 手动确认channel.basicAck(deliveryTag,false);}catch(Exceptione){// 消费失败,拒绝并重新入队channel.basicNack(deliveryTag,false,true);}}}短信消费者结构同理,监听SMS_QUEUE即可。一个交换机绑定了两个队列,同一条消息会同时投递到邮件队列和短信队列,实现并行处理。
七、消息可靠性保障
消息从生产到消费要经过多个环节,任何一个环节都可能丢消息。
7.1 生产者确认机制
publisher-confirm-type:correlated# 异步确认,性能好rabbitTemplate.setConfirmCallback((correlationData,ack,cause)->{if(!ack){System.err.println("消息未到达Exchange,原因:"+cause);// 记录日志,重发等处理}});7.2 消费者手动ACK
默认是自动确认(auto),消息一拿到就标记消费成功,但如果业务代码报异常,消息就丢了。改为手动确认(manual),业务成功后调basicAck,失败调basicNack。
八、死信队列
消息变成"死信"的三种情况:
- 消息被消费者reject(basicReject/basicNack)且不重新入队
- 消息TTL过期(队列或消息设置了过期时间)
- 队列达到最大长度,新消息被挤出去
死信队列的配置思路:给正常队列绑定一个死信交换机(DLX),消息变成死信后自动转发到DLX,再由DLX路由到死信队列。
@BeanpublicQueuenormalQueue(){Map<String,Object>args=newHashMap<>();args.put("x-message-ttl",60000);// 消息60秒过期args.put("x-dead-letter-exchange","dlx.exchange");args.put("x-dead-letter-routing-key","dlx.routing.key");returnnewQueue("normal.queue",true,false,false,args);}死信队列常用于:延迟任务(消息过期→死信→消费)、失败消息重试、订单超时取消等场景。
九、常见问题与解决方案
9.1 消息重复消费(幂等性)
网络抖动导致ACK没及时到达,RabbitMQ会重投消息,消费者就重复处理了。解决方案:
- 业务唯一键校验:消费前查数据库/Redis,已处理则直接ACK跳过
- 乐观锁:update语句加
where status = 0条件 - Redis分布式锁:setnx保证同一消息只处理一次
publicvoidreceive(Stringmessage){StringmsgId=extractMsgId(message);// Redis标记,已处理则跳过BooleanisNew=redisTemplate.opsForValue().setIfAbsent("msg:processed:"+msgId,"1",24,TimeUnit.HOURS);if(Boolean.FALSE.equals(isNew)){return;// 已处理过}// 正常消费逻辑}9.2 消息积压处理
消费速度跟不上生产速度,队列堆积越来越多的消息。应对策略:
- 临时扩容消费者:增加消费者实例数量
- 批量消费:一个消费者一次拉取多条消息处理
- 消息转存:紧急将积压消息转存到另一个队列,后续慢慢消费
- 根因排查:消费者是不是有慢查询?是不是依赖的外部服务超时了?
RabbitMQ用好了就是系统稳定性的护城河,用不好就是给自己挖坑。把可靠性保障和幂等性设计到位,消息队列才能真正发挥价值。
编程学习
技术分享
实战经验