公司动态
RabbitMQ消息队列:异步解耦与业务削峰
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整合RabbitMQ4.1 引入依赖dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-amqp/artifactId/dependency4.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_QUEUEemail.queue;publicstaticfinalStringSMS_QUEUEsms.queue;publicstaticfinalStringORDER_EXCHANGEorder.exchange;publicstaticfinalStringORDER_ROUTING_KEYorder.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);}}五、发送消息RabbitTemplateServicepublicclassOrderService{AutowiredprivateRabbitTemplaterabbitTemplate;publicvoidcreateOrder(OrderDTOorderDTO){// 1. 保存订单数据库操作省略// ...// 2. 异步发送通知消息StringmsgJSON.toJSONString(orderDTO);rabbitTemplate.convertAndSend(RabbitMQConfig.ORDER_EXCHANGE,RabbitMQConfig.ORDER_ROUTING_KEY,msg);// 3. 直接返回不等邮件/短信发送完成return;}}六、接收消息RabbitListenerComponentpublicclassEmailConsumer{RabbitListener(queuesRabbitMQConfig.EMAIL_QUEUE)RabbitHandlerpublicvoidreceive(Stringmessage,Channelchannel,MessagemessageObj)throwsIOException{longdeliveryTagmessageObj.getMessageProperties().getDeliveryTag();try{OrderDTOorderJSON.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。八、死信队列消息变成死信的三种情况消息被消费者rejectbasicReject/basicNack且不重新入队消息TTL过期队列或消息设置了过期时间队列达到最大长度新消息被挤出去死信队列的配置思路给正常队列绑定一个死信交换机DLX消息变成死信后自动转发到DLX再由DLX路由到死信队列。BeanpublicQueuenormalQueue(){MapString,ObjectargsnewHashMap();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){StringmsgIdextractMsgId(message);// Redis标记已处理则跳过BooleanisNewredisTemplate.opsForValue().setIfAbsent(msg:processed:msgId,1,24,TimeUnit.HOURS);if(Boolean.FALSE.equals(isNew)){return;// 已处理过}// 正常消费逻辑}9.2 消息积压处理消费速度跟不上生产速度队列堆积越来越多的消息。应对策略临时扩容消费者增加消费者实例数量批量消费一个消费者一次拉取多条消息处理消息转存紧急将积压消息转存到另一个队列后续慢慢消费根因排查消费者是不是有慢查询是不是依赖的外部服务超时了RabbitMQ用好了就是系统稳定性的护城河用不好就是给自己挖坑。把可靠性保障和幂等性设计到位消息队列才能真正发挥价值。