拓冰建站拓冰建站
首页 / 资讯中心 / 正文

RabbitMQ实战:Spring Boot集成、死信队列与广播模式全解析

简介RabbitMQ代码案例是一套面向Java后端开发者的消息队列实战资源聚焦于分布式系统异步通信、应用解耦与削峰填谷等典型问题。压缩包共299个文件约396KB包含39个Java源码、39个class文件、大量xml配置及properties、jsp等辅助资源结构清晰便于导入IDE直接运行和修改。案例覆盖RabbitMQ核心概念与API调用包括ConnectionFactory连接创建、Channel通道使用、queueDeclare声明队列、basicPublish发送消息、basicConsume消费回调并演示direct、topic、fanout等交换机类型的路由规则同时涉及生产者确认、消费者手动ACK、死信队列、消息TTL等可靠性机制。已有8098人学习该资源说明内容贴合实际开发需求。通过阅读和运行示例读者能快速掌握生产者-消费者模型构建方法理解消息如何从交换机路由到队列学会配置连接、处理异常及保证消息不丢失整套案例是一份轻量而实用的RabbitMQ入门到进阶参考。 做了几年Java后端RabbitMQ算是我用得最多的消息中间件。从最初的一个内部报表系统要用定时任务轮询数据库到后来整套订单体系全部改成MQ异步解耦性能和稳定性完全是两码事。这篇文章我把RabbitMQ从零到一的项目代码案例完整梳理一遍包含Spring Boot集成、广播模式、死信队列、JSON消息传参这些高频场景所有代码都是可以直接复制到项目里跑的版本。不管你是刚接触消息队列的新手还是已经用了一段时间想系统排查一下重复消费、消息积压这些坑的老手这篇都能给你一些参考。1. 写RabbitMQ代码之前先把这几个核心概念吃透1.1 消息中间件到底解决了什么问题很多人一上来就写代码结果连自己为什么要用MQ都说不清楚。面试的时候这个问题也几乎是必问的你为什么在项目里引入RabbitMQ我自己的理解就三条异步、削峰、解耦。异步好理解比如用户下单后要发通知、加积分、更新统计这些串行做要两三秒丢到MQ里异步消费接口直接返回体验完全不一样。削峰是应对瞬时流量比如秒杀场景请求先全部打到MQ后端按自己的消费能力慢慢处理不至于把数据库打崩。解耦是让生产者和消费者互不感知订单系统只管发消息至于谁消费、什么时候消费订单系统完全不关心。1.2 核心模型生产者、消费者、队列、交换机RabbitMQ的模型说起来就四样东西生产者Producer发消息交换机Exchange收消息并按规则路由队列Queue存消息消费者Consumer从队列里取消息处理。新手最容易绕晕的是交换机和队列的关系。记住一句话消息不是直接进队列的而是先到交换机再由交换机根据RoutingKey和绑定规则投递到对应队列。交换机有四种类型Direct精确匹配、Fanout广播给所有绑定队列、Topic通配符匹配、Headers按消息头匹配。实际项目中90%的场景用Direct和Fanout就够了Topic用来做复杂路由。注意生产者只跟交换机打交道消费者只跟队列打交道。这个认知能帮你避免很多配置上的混乱。1.3 代码案例的技术栈选型我下面的代码案例基于Spring Boot 2.7.x Spring AMQP这是目前Java生态里最主流的组合。Spring Boot对RabbitMQ做了很好的封装RabbitTemplate负责发送RabbitListener注解负责消费绝大多数场景不用直接操作底层的Channel API。如果你用的是C#/.NET原理完全一样只是客户端库从Spring AMQP换成了RabbitMQ.ClientExchange、Queue、RoutingKey这套概念是通用的。文末我会顺带提一下C#的对接要点。2. 环境准备Docker安装和Windows本地部署2.1 Docker一键启动最省事的方式开发环境用Docker装RabbitMQ是效率最高的方式不用处理Erlang依赖这些破事。注意一定要带上management标签否则没有Web管理界面。docker run -d --name rabbitmq \ -p 5672:5672 -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASSadmin123 \ rabbitmq:3.12-management5672是AMQP协议端口客户端连接走这个15672是Web管理界面端口浏览器访问http://localhost:15672就能看到控制台。RABBITMQ_DEFAULT_USER和RABBITMQ_DEFAULT_PASS是初始化时自动创建的账号省得再手动配。2.2 Windows本地安装的注意点Windows下安装RabbitMQ稍微折腾一些。先去官网下载Erlang和RabbitMQ安装包版本之间有对应关系必须匹配否则服务起不来。安装完RabbitMQ后在安装目录的sbin文件夹下打开命令行执行rabbitmq-plugins enable rabbitmq_management这一步是启用管理插件装完重启服务浏览器访问http://localhost:15672。默认账号是guest/guest但要注意guest账号默认只能在localhost登录如果要从远程连必须另外创建用户并授权。2.3 管理控制台怎么用控制台里我最常用的几个页面Queues页面看队列的消息积压数量Exchanges页面看交换机绑定关系Connections页面看客户端连接状态。排查问题第一步永远是先上控制台看一眼确认消息到底有没有进队列。实操心得生产环境建议把Management插件只开放给运维网段因为这个控制台的权限很大能直接删除队列、清空消息被误操作一次代价极高。3. Spring Boot集成RabbitMQ一套完整可跑的代码案例3.1 引入依赖和基础配置dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency配置文件application.ymlspring: rabbitmq: host: localhost port: 5672 username: admin password: admin123 virtual-host: / publisher-confirm-type: correlated publisher-returns: true listener: simple: acknowledge-mode: manual prefetch: 5publisher-confirm-type: correlated开启发送方确认回调publisher-returns: true开启消息不可达时的退回机制这两项是保证消息可靠性的关键配置。acknowledge-mode: manual表示消费者手动确认prefetch: 5限制消费者预取数量避免某个消费者一次性拿走太多消息导致负载不均。3.2 硬编码队列和交换机还是用配置类我推荐用Java配置类统一声明队列、交换机、绑定关系这样代码即文档项目里所有消息路由一目了然。看一个标准的Direct模式配置Configuration public class RabbitDirectConfig { public static final String DIRECT_QUEUE direct.queue; public static final String DIRECT_EXCHANGE direct.exchange; public static final String DIRECT_ROUTING_KEY direct.routing.key; Bean public Queue directQueue() { return QueueBuilder.durable(DIRECT_QUEUE).build(); } Bean public DirectExchange directExchange() { return new DirectExchange(DIRECT_EXCHANGE, true, false); } Bean public Binding directBinding() { return BindingBuilder.bind(directQueue()) .to(directExchange()) .with(DIRECT_ROUTING_KEY); } }durable表示队列持久化重启不丢。DirectExchange三个参数分别是名称、是否持久化、是否自动删除生产环境第一个参数用true第二个用false。这里有个细节如果队列参数和声明时不一致比如持久化状态、死信参数启动时会直接报错所以队列参数一旦定下来后续尽量不要改。3.3 生产者用RabbitTemplate发消息Service public class OrderMessageProducer { Autowired private RabbitTemplate rabbitTemplate; Resource private RabbitTemplate confirmRabbitTemplate; PostConstruct public void init() { // 发送确认回调 confirmRabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (!ack) { log.error(消息发送失败: {}, cause: {}, correlationData, cause); } }); } public void sendOrderMessage(OrderMessage message) { CorrelationData correlationData new CorrelationData(UUID.randomUUID().toString()); confirmRabbitTemplate.convertAndSend( RabbitDirectConfig.DIRECT_EXCHANGE, RabbitDirectConfig.DIRECT_ROUTING_KEY, message, correlationData ); } }convertAndSend方法会自动把对象序列化成JSON字节流。CorrelationData携带业务唯一的消息ID回调里可以通过它关联到具体是哪条消息发送失败。我在项目里会把发送失败的消息落库配合定时任务做补偿重发这是保证消息不丢的最朴素也最可靠的办法。3.4 消费者RabbitListener监听队列Component public class OrderMessageConsumer { RabbitListener(queues RabbitDirectConfig.DIRECT_QUEUE) public void onMessage(OrderMessage message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws Exception { try { // 业务处理 orderService.process(message); // 手动确认 channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(消费失败, e); // 第三个参数requeuefalse消息进入死信队列而不是无限重试 channel.basicNack(deliveryTag, false, false); } } }手动确认模式下消费逻辑必须包在try-catch里。basicAck确认成功basicNack拒绝消息。第三个参数requeue要重点说明如果传true消息会回到原队列头部继续尝试如果业务逻辑一直报错消息就会无限循环把消费线程彻底卡死。正确做法是requeuefalse让消息转到死信队列由单独的死信消费者做补偿或人工介入。3.5 JSON消息的坑序列化和反序列化Spring Boot默认的消息转换器是SimpleMessageConverter只能处理字符串和字节数组。如果把对象直接塞给RabbitTemplate要么报错要么序列化成Java原生格式别的服务消费不了。统一改成JSON转换器Bean public MessageConverter jsonMessageConverter() { return new Jackson2JsonMessageConverter(); }在配置类里加上这个Bean后生产者发送对象会自动转JSON消费者接收时RabbitListener的方法参数如果是自定义对象也会自动从JSON反序列化。这里有一个非常容易踩的坑消费者反序列化报ClassNotFoundException。原因是生产者端把对象的全限定类名写进了消息头消费者项目里如果类名或包路径不一致就会报错。解决方法是传输对象单独建一个公共模块或者用Map/JsonNode接收后再手动转。4. 五种消息模型不同场景用哪个代码怎么写4.1 Work Queue工作队列能者多劳Work模式是最简单也最常见的一个队列多个消费者竞争消费。Spring Boot里只需要在多个类里用RabbitListener监听同一个队列即可配合prefetch参数实现能者多劳。RabbitListener(queues work.queue) public void consumerA(String message) { Thread.sleep(1000); } RabbitListener(queues work.queue) public void consumerB(String message) { Thread.sleep(3000); }关键是prefetch默认值是250这个值是每个消费者预取的最大未确认消息数。如果A消费者处理快、B处理慢默认配置下两个消费者会均分消息而不是按处理能力分配。把prefetch设成1每次只取一条消息处理完确认后再取下一条才能真正实现能者多劳。我在订单处理场景用的prefetch5在业务逻辑简单的日志推送场景用的prefetch50要根据单条消息处理耗时去调没有固定值。4.2 Fanout广播模式一次发送所有队列收到Fanout交换机不关心RoutingKey它会把消息复制投递给所有绑定的队列适合发布公告、刷新缓存、同步全量配置这类场景。看代码Configuration public class RabbitFanoutConfig { public static final String FANOUT_EXCHANGE fanout.exchange; Bean public Queue queueA() { return QueueBuilder.durable(fanout.queue.A).build(); } Bean public Queue queueB() { return QueueBuilder.durable(fanout.queue.B).build(); } Bean public FanoutExchange fanoutExchange() { return new FanoutExchange(FANOUT_EXCHANGE); } Bean public Binding bindingA() { return BindingBuilder.bind(queueA()).to(fanoutExchange()); } Bean public Binding bindingB() { return BindingBuilder.bind(queueB()).to(fanoutExchange()); } }这个模式在Ruoyi这类后台管理框架里很常用比如用户修改了角色权限需要同时通知多个微服务刷新本地权限缓存。生产者发一条消息所有相关服务都能收到不用单独维护每个服务的调用关系。4.3 Topic主题模式通配符做灵活路由Topic模式下RoutingKey用.分隔单词交换机支持通配符*匹配一个单词#匹配零个或多个单词。比如订单服务产生一条订单创建事件路由键是order.created同时绑定order.*和order.#的队列都能收到。这个模式适合做事件驱动架构。Bean public TopicExchange topicExchange() { return new TopicExchange(topic.exchange); } Bean public Queue orderAllQueue() { return QueueBuilder.durable(order.all.queue).build(); } Bean public Queue orderCreatedQueue() { return QueueBuilder.durable(order.created.queue).build(); } Bean public Binding bindingAll() { return BindingBuilder.bind(orderAllQueue()) .to(topicExchange()) .with(order.#); } Bean public Binding bindingCreated() { return BindingBuilder.bind(orderCreatedQueue()) .to(topicExchange()) .with(order.created); }我建议路由键的命名规范统一用业务模块.事件类型比如order.paid、order.cancelled、user.registered这样通过命名就能看出消息的业务含义后续加消费者只需要新增绑定对已有链路毫无影响。4.4 RPC模式MQ做同步调用能不用就别用RabbitMQ支持RPC模式客户端发送消息时带上replyTo指定回调队列服务端处理完把结果发回回调队列客户端阻塞等待。这听起来很酷但我个人的建议是如果是同步调用直接用HTTP或Feign更简单。MQ的定位是异步解耦做同步RPC等于把简单问题复杂化还要处理响应超时、回调队列堆积等额外问题。5. 死信队列消息兜底方案和可靠性保障5.1 什么情况消息会进死信队列死信队列是RabbitMQ的兜底机制。三种情况消息会进入死信消费者basicNack且requeuefalse、消息TTL过期、队列达到最大长度。死信的价值在于业务上的失败消息不会丢而是被隔离到一个专门的队列由独立的消费者做补偿处理。比如订单支付回调消费失败三次正常消费线程退出重试消息进死信队列死信消费者收到后做告警和人工干预。5.2 死信队列配置代码案例Configuration public class RabbitDlxConfig { public static final String DEAD_QUEUE order.dead.queue; public static final String DEAD_EXCHANGE order.dead.exchange; public static final String DEAD_ROUTING_KEY order.dead.routing.key; Bean public Queue orderQueue() { return QueueBuilder.durable(order.queue) // 声明死信交换机 .deadLetterExchange(DEAD_EXCHANGE) // 声明死信路由键 .deadLetterRoutingKey(DEAD_ROUTING_KEY) .build(); } Bean public Queue deadQueue() { return QueueBuilder.durable(DEAD_QUEUE).build(); } Bean public DirectExchange deadExchange() { return new DirectExchange(DEAD_EXCHANGE); } Bean public Binding deadBinding() { return BindingBuilder.bind(deadQueue()) .to(deadExchange()) .with(DEAD_ROUTING_KEY); } }deadLetterExchange和deadLetterRoutingKey这两个参数是加在业务队列上的作用是告诉RabbitMQ当这条队列里的消息变成死信时转发到哪个交换机。这个配置特别容易漏很多人只建了业务队列忘了给业务队列添加死信参数导致消息被basicNack后直接丢弃。5.3 死信消息的消费和补偿Component public class DeadLetterConsumer { RabbitListener(queues RabbitDlxConfig.DEAD_QUEUE) public void onDeadMessage(OrderMessage message) { log.error(收到死信消息业务单号: {}, 补偿处理, message.getOrderId()); // 1. 查数据库确认业务状态 // 2. 如果业务未完成调用补偿接口 // 3. 如果补偿也失败发送告警通知人工处理 } }这里有个很容易犯的错误死信消费者如果处理失败又basicNack消息会再次进入死信队列形成死循环。我的做法是死信消费者固定用自动确认模式消费失败就把消息详细信息落库同时发送钉钉/企业微信告警由人工介入。宁可消息重复处理也不要让它无限循环。5.4 消息积压问题30分钟能压多少消息热词里有一条rabbitmq 死信30分种会压多少这其实问的是消息积压的量级估算。我实测过一组数据供参考单条消息平均1KB消费者单线程处理耗时50ms30分钟能消费约3600条30×60×1000/50。如果单条消息4KB且处理耗时200ms30分钟只能消费9000条。所以处理积压问题的核心思路是先让消费者把消息快速确认掉哪怕先落库再异步慢慢处理业务逻辑。把耗时的下游调用从消费线程里拆出去消费能力可以提升几十倍。我在一个数据同步项目里就是这么做的消费瓶颈从每秒80条提升到每秒2000条靠的就是消费即确认、业务延后处理这个思路。6. 常见问题和排查技巧实录6.1 消费者收不到消息先从这几个地方查这是群里被问烂了的问题。我的排查顺序先上管理控制台的Queues页面看消息是否存在、处于Ready还是Unacked状态再看消费者的连接状态是否在线然后用控制台的Get Message功能手动取一条看看格式。如果控制台显示消息在队列里但消费者一直不收百分之八九十是绑定的RoutingKey对不上。消费者监听的是A队列生产者的routing key把消息路由到了B队列这在控制台的Exchange页面看绑定关系一目了然。6.2 重复消费问题消息队列不保证不重复只保证不丢失。网络抖动时消费者处理完消息准备确认但连接断了消息会重新投递这就是重复消费的来源。解决思路就一条消费逻辑做到幂等。最简单的方式是消息带上唯一业务ID消费前查重。public void process(OrderMessage message) { // 利用数据库唯一索引做幂等 int count orderProcessLogMapper.insertIfAbsent(message.getOrderId()); if (count 0) { // 已处理过直接返回 return; } // 正常业务处理 orderService.process(message); }6.3 连接频繁断开如果客户端连接上了但每隔一段时间自动断开优先检查心跳超时配置。RabbitMQ默认心跳60秒客户端如果没有及时发送心跳包服务端会判定连接失效。在配置里显式设置spring: rabbitmq: requested-heartbeat: 30 connection-timeout: 15000另外注意防火墙和负载均衡设备的空闲连接回收时间如果它比MQ的心跳时间短连接会被中间设备切掉这种情况要调整心跳间隔适配网络环境。6.4 生产环境参数调优表我总结了一份常用的参数建议表可以按项目实际情况调整参数默认值建议值说明prefetch2501-50单消费者预取消息数根据单条处理耗时调整消费者线程数CPU核数根据队列QPS调整通过concurrency和max-concurrency设置消息TTL无按业务需要设置后过期消息自动进死信队列队列最大长度无限制按业务评估超过后最老的消息进死信或丢弃确认模式automanual重要业务必须手动确认6.5 C#对接RabbitMQ的几个要点热词里C#相关的搜索量不小我简单提一下。C#使用RabbitMQ.Client库核心代码和Java差不多ConnectionFactory创建连接IModel上声明队列和交换机BasicPublish发布消息EventingBasicConsumer接收消息。JSON序列化在C#里用JsonSerializer或者Newtonsoft.Json手动转成字符串再发送。C#端的坑主要集中在连接工厂的AutomaticRecoveryEnabled要设为true否则网络抖动后连接不会自动恢复。我自己在实际项目里最深刻的一个体会是RabbitMQ部署和基础用法都不难难的是消息可靠性的整体设计。交换机、队列、死信、幂等、补偿机制这五件事在项目初期就规划好后面能省掉无数排查问题的时间。你从最简单的Direct模式开始跑通一条消息再逐步引入广播和死信循序渐进很快就能摸清楚它的脾性。最后再分享一个小建议生产环境务必开启publisher-confirm和手动ack这两项配置是消息不丢的底线别为了省事用默认配置。本文还有配套的精品资源点击获取
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门