
简介这是一套面向Java开发者的RabbitMQ入门到进阶实战代码案例。案例围绕AMQP协议从生产者、消费者、交换机、队列等核心概念出发演示了基于amqp-client库的消息发送与接收流程并扩展到direct、topic、fanout等多种交换机路由策略以及事务、发布确认、死信队列、消息TTL等可靠性机制。压缩包大小仅396KB共包含299个文件其中39个Java源码与39个class可直接对应学习另有188个XML配置和properties资源文件用于工程搭建辅以少量JSP页面及Maven相关文件便于在IDEA等环境中运行调试。已有8098人学习下载适合需要快速掌握消息中间件实际用法、或正在设计分布式异步系统的开发人员参考。通过阅读和运行这些示例可以理解RabbitMQ API的调用细节与消息路由原理少走弯路并快速落地到自己的项目中。 搞消息队列这几年RabbitMQ 是我用得最顺手的一个。很多人一提到 RabbitMQ 就只会发个 hello world真到了生产环境广播、死信、JSON 序列化、多语言对接全得靠代码案例一点点趟坑。这篇就用我实际跑通过的项目代码把 RabbitMQ 从安装、启动、SpringBoot 集成、C# 推送、死信队列到 Docker 集群部署完整串一遍。适合刚接触消息队列的开发者也适合已经在用但想补全细节的同行。1. 项目背景与整体设计思路1.1 为什么选 RabbitMQ而不是其他消息队列做技术选型时很多人爱纠结 Kafka 还是 RocketMQ但 RabbitMQ 在中小型系统、企业内部服务间异步通信上优势非常明显部署轻量、管理界面友好、路由规则灵活、多语言客户端成熟。更重要的是它的 AMQP 协议模型交换机、队列、绑定几乎覆盖了你日常能想到的所有消息场景。我这边有个真实业务前端提交工单后系统需要同时通知多个服务去处理比如发短信、写日志、更新统计。如果同步调用一个服务挂了整个流程就卡住。用 RabbitMQ 做广播 fanout 模式消息发到交换机所有绑定的队列都会收到互不影响。另一个场景是订单超时未支付需要自动关闭这就要用到死信队列加延迟路由。1.2 业务场景梳理广播、异步、死信分别解决什么问题场景一广播通知。一条消息要让所有监听方都收到用 fanout 交换机最合适。场景二异步解耦。用户注册成功后只需要写库邮件验证码之类的异步发出去用户无感知。场景三死信兜底。消息消费失败、超时未被确认进入死信队列后续做补偿处理或人工介入。这三个场景对应了 RabbitMQ 最常见的三类用法也正好是面试中常被追问的“你在项目中哪里用到了 MQ”。能讲清楚场景比背概念有用得多。1.3 环境准备Windows 安装、Docker 部署、集群搭建环境是第一步也是最容易踩坑的一步。Windows 下安装 RabbitMQ 需要先装 Erlang而且版本要匹配。我之前试过 Erlang 24 配 RabbitMQ 3.9 没问题后来换机器装成 Erlang 26 就出现启动后管理界面打不开的情况。官方有兼容性对照表装之前一定要看一眼。Docker 方式相对省心一条命令就能跑起来docker run -d --name rabbitmq \ -p 5672:5672 -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASSadmin123 \ rabbitmq:3.12-management这样装完自带管理插件浏览器访问 15672 端口就能看到界面。如果你要搭集群记住几个关键点节点 cookie 必须一致、每个节点的 hostname 不能被解析成 127.0.0.1、用rabbitmqctl join_cluster把新节点加入集群。Docker 下做集群我建议用--network host或者 docker-compose 固定容器名否则节点间通讯经常会出幺蛾子。2. 核心代码案例一SpringBoot 集成 RabbitMQ广播 JSON2.1 引入依赖与基础配置SpringBoot 集成 RabbitMQ 算是 Java 生态里最省事的方案。首先在pom.xml里加依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency然后在application.yml里配置连接信息spring: rabbitmq: host: localhost port: 5672 username: admin password: admin123 virtual-host: / publisher-confirm-type: correlated publisher-returns: true这里有两个容易被忽略的参数publisher-confirm-type和publisher-returns。前者开启发送端确认消息到达交换机后回调告诉你结果后者开启消息路由不到队列时的返回回调。生产环境一定要开不然消息丢失了你根本不知道。2.2 定义交换机、队列和绑定关系RabbitMQ 的核心是交换机、队列、绑定三者之间的关系。在我的广播案例里项目里通常用注解方式直接声明Configuration public class RabbitConfig { public static final String FANOUT_EXCHANGE exchange.notify; public static final String QUEUE_SMS queue.sms; public static final String QUEUE_LOG queue.log; Bean public FanoutExchange fanoutExchange() { return new FanoutExchange(FANOUT_EXCHANGE); } Bean public Queue smsQueue() { return QueueBuilder.durable(QUEUE_SMS).build(); } Bean public Queue logQueue() { return QueueBuilder.durable(QUEUE_LOG).build(); } Bean public Binding smsBinding() { return BindingBuilder.bind(smsQueue()).to(fanoutExchange()); } Bean public Binding logBinding() { return BindingBuilder.bind(logQueue()).to(fanoutExchange()); } }fanout 交换机不关心 routingKey只要消息发到交换机上所有绑定的队列都能收到。如果你用 direct 或 topic那就要在 bind 的时候指定 routingKey发送时也得带对应 key。曾经有个同事把 direct 当 fanout 用routingKey 写错导致消息全部进入黑洞排查了一下午其实就是没理解模型。2.3 生产者推送 JSON 消息把 JSON 放入 RabbitMQ是实际项目里最常见的操作。SpringBoot 里可以手动转 JSON也可以配置消息转换器让convertAndSend自动把对象序列化成 JSON。我推荐配置一个全局 Jackson 转换器Bean public MessageConverter messageConverter() { return new Jackson2JsonMessageConverter(); }生产者推送Service public class NotifyService { Autowired private RabbitTemplate rabbitTemplate; public void sendNotify(NotifyMessage message) { CorrelationData correlationData new CorrelationData(UUID.randomUUID().toString()); rabbitTemplate.convertAndSend( RabbitConfig.FANOUT_EXCHANGE, , message, correlationData ); } }NotifyMessage就是一个普通 POJO字段包括id、content、createTime。加了Jackson2JsonMessageConverter后RabbitTemplate 会把对象转成 JSON 字符串再包装成消息体发送。注意如果你在 fanout 模式下把 routingKey 随便填一个字符串消息依然能发出去因为 fanout 根本不检查 key这也容易造成“明明绑定了队列却收不到消息”的错觉。2.4 消费者监听与消息确认消费者就更简单了用RabbitListener注解就能监听队列Component public class NotifyConsumer { RabbitListener(queues RabbitConfig.QUEUE_SMS) public void handleSms(NotifyMessage message) { System.out.println(发送短信: message.getContent()); } RabbitListener(queues RabbitConfig.QUEUE_LOG) public void handleLog(NotifyMessage message) { System.out.println(写日志: message.getContent()); } }这里有个非常重要的参数RabbitListener(queues ..., ackMode MANUAL)。默认情况下 SpringBoot 的监听器是自动确认只要方法执行完就 ack。但如果你在处理消息时抛了异常Spring 会自动 requeue导致消息一直在队列头部打转形成死循环。常见做法是改成手动确认RabbitListener(queues RabbitConfig.QUEUE_SMS, ackMode MANUAL) public void handleSms(NotifyMessage message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException { try { System.out.println(发送短信: message.getContent()); channel.basicAck(tag, false); } catch (Exception e) { channel.basicNack(tag, false, true); } }basicNack的第三个参数是requeue设为 true 表示重新放回队列。如果确认消息确实有问题建议设为 false配合死信队列做兜底而不是无限重试。3. 核心代码案例二C# 推送 RabbitMQ3.1 C# 端连接与发送很多公司是 Java 后端 C# 桌面端或 .NET 服务并存这时候 RabbitMQ 跨语言优势就体现出来了。C# 用官方 RabbitMQ.Client 库NuGet 里直接装Install-Package RabbitMQ.Client发送端核心代码using RabbitMQ.Client; using System.Text; using System.Text.Json; var factory new ConnectionFactory() { HostName localhost, Port 5672, UserName admin, Password admin123 }; using var connection factory.CreateConnection(); using var channel connection.CreateModel(); var notify new { id Guid.NewGuid().ToString(), content hello from csharp, createTime DateTime.Now }; string json JsonSerializer.Serialize(notify); byte[] body Encoding.UTF8.GetBytes(json); channel.BasicPublish( exchange: exchange.notify, routingKey: , basicProperties: null, body: body );注意C# 端声明交换机的时候不能用ExchangeDeclarePassive去检查但不存在就报错。通常做法是确保 C# 里也有一样的交换机声明逻辑或者由 Java 服务统一把交换机队列建好C# 只负责发送。我在实际项目里是让 Java 服务管理所有交换机队列声明C# 只推送避免两边声明不一致。3.2 让 JSON 消息被多语言消费者正确消费跨语言消费有个细节消息体里的字段命名风格。Java 默认字段是驼峰createTimeC# 序列化默认也是驼峰但如果 C# 侧反序列化成CreateTime属性System.Text.Json默认区分大小写就会反序列化失败。解决方法是统一字段命名或在 C# 反序列化时配置PropertyNameCaseInsensitive truevar options new JsonSerializerOptions { PropertyNameCaseInsensitive true }; var msg JsonSerializer.DeserializeNotifyMessage(json, options);还有一点RabbitMQ 的消息头里content_type一定要是application/json。Java 的Jackson2JsonMessageConverter会自动设置但 C# 手动BasicPublish的时候如果不设置basicProperties默认是text/plain。虽然消费者按 JSON 解也能成功但有些框架会根据 content_type 做消息转换建议显式设置var props channel.CreateBasicProperties(); props.ContentType application/json; channel.BasicPublish(exchange.notify, , props, body);3.3 C# 消费端的常用写法C# 消费端用 EventingBasicConsumer 订阅using var channel connection.CreateModel(); channel.QueueDeclare(queue.result, durable: true, exclusive: false, autoDelete: false); var consumer new EventingBasicConsumer(channel); consumer.Received (model, ea) { string json Encoding.UTF8.GetString(ea.Body.ToArray()); Console.WriteLine($收到消息: {json}); channel.BasicAck(ea.DeliveryTag, false); }; channel.BasicConsume(queue: queue.result, autoAck: false, consumer: consumer);这里我特别强调autoAck: false。C# 消费端如果设成 true消息一收到就确认万一处理逻辑崩了消息就丢了。设成 false 配合手动BasicAck至少在异常时能通过BasicNack把消息放回队列或进死信。4. 死信队列实战延迟消息与异常兜底4.1 死信队列的原理死信队列Dead Letter Queue本质上就是普通的 RabbitMQ 队列只不过它专门接收“没有正常被消费”的消息。消息变成死信的条件有三个消费者basicNack且requeuefalse、消息过期未被消费、队列长度达到上限。利用这个机制可以模拟延迟队列给普通队列设置消息过期时间 TTL消息过期后自动转到死信交换机死信交换机再把消息路由到真正的消费队列。RabbitMQ 本身没有延迟队列插件时这是最经典的延迟消息实现方案。4.2 代码实现普通队列绑定死信交换机SpringBoot 里声明一个带死信属性的队列Configuration public class DelayConfig { public static final String DELAY_EXCHANGE exchange.delay; public static final String DELAY_QUEUE queue.delay; public static final String DEAD_EXCHANGE exchange.dead; public static final String DEAD_QUEUE queue.dead; Bean public DirectExchange delayExchange() { return new DirectExchange(DELAY_EXCHANGE); } Bean public DirectExchange deadExchange() { return new DirectExchange(DEAD_EXCHANGE); } Bean public Queue delayQueue() { return QueueBuilder.durable(DELAY_QUEUE) .deadLetterExchange(DEAD_EXCHANGE) .deadLetterRoutingKey(dead) .build(); } Bean public Queue deadQueue() { return QueueBuilder.durable(DEAD_QUEUE).build(); } Bean public Binding delayBinding() { return BindingBuilder.bind(delayQueue()).to(delayExchange()).with(delay); } Bean public Binding deadBinding() { return BindingBuilder.bind(deadQueue()).to(deadExchange()).with(dead); } }生产者设置消息过期时间MessagePostProcessor processor message - { message.getMessageProperties().setExpiration(1800000); return message; }; rabbitTemplate.convertAndSend(DELAY_EXCHANGE, delay, order, processor);setExpiration的单位是毫秒1800000 就是 30 分钟。这条消息发到 delay 队列后不会被立即消费而是要等 30 分钟过期然后转进死信交换机最终到达真正的消费队列。4.3 30分钟积压场景分析与参数调优热词里有“rabbitmq 死信30分种会压多少”我觉得这个问题的本质是延迟消息设为 30 分钟后同时会有多少消息积压在延迟队列里。要估算积压量核心就两个指标每秒进入延迟队列的消息数和延迟时间。比如系统每秒产生 100 条超时订单延迟 30 分钟1800 秒那累积在队列里的消息数大约就是 100 * 1800 180000 条。这是理论积压量。这数字影响什么影响队列内存、磁盘持久化、以及消息过期后瞬间转入死信队列的冲击力。如果 18 万条消息同一时间点过期那死信队列会瞬间收到大量消息消费者可能扛不住。解决办法是把过期时间做离散化比如在业务上将延迟时间增加随机偏移量让过期时间分散。另一个办法是给死信队列配置多个消费者加并发或者用 RabbitMQ 的x-max-priority来控制某些紧急消息优先消费。实际测试中普通单机 RabbitMQ 几万条消息没问题但 18 万条消息同时触发流转管理界面会明显卡顿消费者若不加prefetch限制内存占用直接飙升。所以我建议给消费者设置prefetchCount用 SimpleMessageListenerContainer 的时候Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(ConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); factory.setPrefetchCount(100); factory.setAcknowledgeMode(AcknowledgeMode.MANUAL); return factory; }5. 常见问题与排查技巧5.1 启动失败、端口占用与管理界面打不开Windows 下 RabbitMQ 启动失败八成都出在 Erlang 版本不匹配或主机名解析上。解决办法用管理员身份打开命令行先运行rabbitmq-service stop再rabbitmq-service uninstall清掉 Erlang 的 cookie 缓存重装匹配版本。端口方面5672 是 AMQP 端口15672 是管理界面端口。用 Docker 时如果 15672 被占用改映射端口即可但 5672 必须跟客户端配置保持一致。排查命令rabbitmqctl status rabbitmq-diagnostics -q ping rabbitmq-plugins enable rabbitmq_management5.2 消息堆积、重复消费与确认机制消息堆积常见原因有两个生产者速度远大于消费者速度、消费者处理逻辑太慢。前者可以增加消费者实例或提高 prefetch后者需要优化业务代码甚至把下游写库改成批量写。还有一种隐蔽情况消费者一直接收消息但每一条都抛异常 requeue看起来队列长度不变实际是在死循环这个问题只能靠日志定位。重复消费是消息系统永远绕不开的话题。RabbitMQ 不保证 exactly once所以消费者必须具备幂等性。最简单的方式是消息里带上业务唯一 ID消费前查一下 Redis 或数据库已经处理过就直接 ack。5.3 Docker 集群部署的常见坑Docker 下搭 RabbitMQ 集群我踩过的坑有几个。第一是容器重启后 hostname 变化导致集群节点失效解决方法是固定 hostname 或用--hostname参数。第二是节点间 cookie 不一致/var/lib/rabbitmq/.erlang.cookie要保证一致。第三是用-p 5672:5672这种端口映射方式做集群时节点间通讯走了随机端口防火墙拦掉就连不上。建议直接用 docker-compose定义好网络和固定容器名避免不必要的麻烦。最后再分享一个小技巧如果你正在准备面试或刚接手 RabbitMQ 项目建议你自己动手搭一套最小的代码案例一个生产者定时推送 JSON 消息两个消费者分别监听其中一个消费失败进死信队列再写一个定时任务扫描死信队列做补偿。这套案例跑通后你对交换机、队列、绑定、确认机制、死信流转的理解会彻底贯通。很多概念看十遍文档不如亲手看一眼消息进到死信队列时管理界面上的流转路径。本文还有配套的精品资源点击获取