ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

RabbitMQ从入门到生产:交换机、死信队列与集群实战

RabbitMQ从入门到生产:交换机、死信队列与集群实战 简介RabbitMQ是流行的开源消息队列中间件采用AMQP协议实现可靠的异步通信。这套代码围绕Java语言展示生产者、消费者、交换机和队列的协作方式帮助读者理解消息代理在分布式系统解耦与削峰中的作用。压缩包内299个文件包含Java源码、XML配置、class字节码、properties配置和少量依赖库与日志文件大小约396KB便于下载后直接查看代码结构和运行验证。案例给出完整的发送与接收程序覆盖连接创建、队列声明、消息发布与回调消费等关键步骤还梳理了主题交换机与直接交换机的路由差异并涉及消息确认、持久化等可靠性机制。已有超过八千人学习适合刚接触消息队列的开发者循序渐进地练习。借助这套代码可较快掌握RabbitMQ核心API与常见设计模式为生产环境消息链路改造打下基础。 前段时间帮朋友排查线上RabbitMQ消息丢失的问题他问了一个几乎所有初学者都会问的点“同一个交换机消息发出去也没报错为什么消费者就是收不到”最后找到原因绑定队列时RoutingKey写错了而且生产者还没开发布确认。这个场景太典型了——很多人能跑通Hello World但一接触到真实业务就露怯。这篇文章我把从安装、核心模型、Spring Boot实战、死信队列到跨语言调用和集群部署的整套代码案例和踩坑记录整理出来适合正在找RabbitMQ案例、准备面试、或者已经在生产环境里被消息问题折磨过的人。1. 装环境那点事Windows与Docker两条路线别在第一步卡太久1.1 Windows安装的隐藏依赖RabbitMQ是用Erlang写的所以Windows下安装必须先装Erlang而且版本必须严格对应。RabbitMQ官方文档每个版本都会列出支持的Erlang版本区间我见过太多人随便装了个新版Erlang结果RabbitMQ服务起不来日志里报Failed to start Erlang之类的错。这里有个经验去RabbitMQ官网的“Install on Windows”页面它会直接给你对应的Erlang下载链接别自己去Erlang官网下最新的。装完之后RabbitMQ默认会注册成Windows服务。启动方式有两种一种是在服务管理器里找到RabbitMQ服务点启动另一种是用命令行管理工具路径在RabbitMQ安装目录的sbin下比如rabbitmq-server.bat start rabbitmqctl status很多人在Windows上碰到“端口被占用”的问题。RabbitMQ默认监听5672AMQP协议端口和15672管理控制台端口。如果你本机装过其他消息中间件或者某些开发工具5672很容易被占。排查命令netstat -ano | findstr 5672如果端口被占要么停掉占用程序要么改RabbitMQ配置。这里我建议优先停占用程序因为改端口后面所有连接配置都要跟着改很烦。1.2 Docker方式我更推荐Windows上装RabbitMQ容易踩坑Docker就省心多了。一条命令拉起来docker run -d --name rabbitmq \ -p 5672:5672 -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASSadmin123 \ rabbitmq:3-management注意官方rabbitmq:3-management镜像已经带了管理插件不需要再手动rabbitmq-plugins enable。如果你拉的是不带management标签的纯rabbitmq:3镜像就算端口映射了15672也打不开控制台。Docker方式还有个好处是换版本方便测试环境想升级就重新拉一个镜像再起容器几秒钟的事。等你看完这篇文章后面的集群章节你会发现Docker部署集群比Windows上手工配Erlang Cookie省心得多。2. 先想清楚再动手交换机、队列、路由键到底在搬什么2.1 三者关系一张图之外的事很多人看RabbitMQ的架构图觉得挺清楚生产者把消息发给交换机交换机根据绑定关系把消息路由到队列消费者从队列拉消息。但一写代码就糊涂原因是没有理解交换机有不同类型。Direct交换机按RoutingKey精确匹配。我发消息时指定的RoutingKey叫order.created队列绑定时的RoutingKey也必须是order.created完全一样才路由。Topic交换机按通配符匹配。order.*能匹配order.createdorder.#能匹配order.created.success这种多级。这种交换机适合做业务分类比如订单消息全部走order.*。Fanout交换机不关心RoutingKey广播给所有绑定的队列。热搜词里有人搜“ruoyi集成springboot集成rabbitmq广播模板”这种场景就适合Fanout比如一个用户登录事件短信服务、日志服务、积分服务各一个队列同时收到消息。我的建议是绝大多数业务用Direct就够了只有明确需要“一对多广播”时才用FanoutTopic适合规则复杂的路由场景。别一上来就追求花哨路由规则越复杂线上排查越痛苦。2.2 消息确认和持久化决定消息会不会丢这是面试必问、生产必踩的坎。我拿真实的丢消息案例来说生产者发消息成功消费者处理失败消息没了。原因就是消费者自动ACK——RabbitMQ默认消费者收到消息就自动确认哪怕处理逻辑抛异常消息也已经从队列里被标记删除了。解决办法是改手动确认channel.basicConsume(queueName, false, consumer); // 第二个参数autoAck设为false // 处理成功后 channel.basicAck(deliveryTag, false); // 处理失败后 channel.basicNack(deliveryTag, false, true); // requeuetrue会重新放回队列再说持久化。光有手动ACK还不够消息默认只存在内存里RabbitMQ一重启就全没了。要做到持久化需要三件套队列声明时durabletrue交换机声明时durabletrue发送消息时设置MessageProperties.PERSISTENT_TEXT_PLAIN或自己构造BasicProperties时deliveryMode2。三个缺一个重启都可能丢消息。有个容易忽略的细节持久化不是把消息同步刷到磁盘的它只是先写内存再异步刷盘。万一RabbitMQ进程被强杀刚发还没落盘的消息照样丢。要做到真正不丢还要生产者开Publisher Confirm就是下面Spring Boot案例里要讲到的publisher-confirm-type。2.3 面试题里的“为什么”面试问RabbitMQ高频问题其实是“为什么用RabbitMQ而不是Kafka”。我的理解是RabbitMQ是面向服务间解耦和复杂路由的追求低延迟、消息可靠投递、精细控制Kafka是面向海量日志和流式处理的追求高吞吐、顺序追加、回放。你做了订单、支付这类强一致性要求的业务RabbitMQ的ACK机制、死信队列、优先级队列这些能力会更顺手。另一个高频题是“RabbitMQ怎么保证消息不丢失”答案就对应上面的生产者确认、队列持久化、消费者手动ACK三件事。这个链条要背熟然后能配合代码讲清楚。3. 从“能跑通”到“敢上生产”一个Spring Boot案例的全部细节3.1 依赖和配置先用Spring Boot集成这是目前Java生态里最主流的用法。引入依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency配置文件里有几个参数是生产级必配的我用注释标注了为什么spring: rabbitmq: host: 127.0.0.1 port: 5672 username: admin password: admin123 publisher-confirm-type: correlated # 开启发布确认发出去没发出去生产者能拿到回调 publisher-returns: true # 消息路由不到队列时回退给生产者 template: mandatory: true # 配合publisher-returns路由失败必须回调 listener: simple: acknowledge-mode: manual # 消费端手动ACK别用自动 prefetch: 10 # 每个消费者每次最多取10条避免积压在客户端prefetch这个参数很容易被忽略。如果不设置RabbitMQ默认是一次性推很多消息给消费者消费者处理不过来内存压力大。设成10就是告诉服务端“我每次最多同时处理10条处理完再给我新的”这是保护消费者的重要手段。3.2 消息发出去不是结束要确认它真的到了定义一个RabbitTemplate。Spring Boot 2.x里publisher-confirm-type: correlated会用CorrelationData作为回调标识Service public class OrderMessageSender { Autowired private RabbitTemplate rabbitTemplate; public void sendOrderCreated(OrderDTO order) { CorrelationData correlationData new CorrelationData(order.getOrderId()); rabbitTemplate.convertAndSend( order.exchange, order.created, order, message - { message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT); return message; }, correlationData ); // convertAndSend是异步发送如果要立刻知道结果可以用waitForConfirms或回调 } }convertAndSend的重载参数里CorrelationData就是用来关联发布确认回调的。你可以单独写一个ConfirmCallbackrabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (!ack) { log.error(消息发送失败: {}, cause: {}, correlationData.getId(), cause); // 这里要做补偿落库标记、重试或者告警 } });注意发布确认回调是异步的生产者代码根本感觉不到网络问题所以线上必须把“发送失败”和“路由失败”的回调接好。我曾见过一个项目没配mandatory和returns结果路由不存在的队列时消息静默丢失查了一天。3.3 消费者的正确打开方式用一个RabbitListener监听队列Component public class OrderCreatedConsumer { RabbitListener(queues order.pending.queue) public void onMessage(OrderDTO order, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { try { // 真正的业务处理比如创建支付单 handleOrder(order); channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(消费失败, e); // 重试三次之后还是失败进死信队列 channel.basicNack(deliveryTag, false, false); // requeue设为false } } }这里有个关键设计basicNack的requeue参数什么时候设true什么时候设false如果是消费接口暂时不可用、网络抖动这种临时性问题设true让消息重新排队过一会再消费。如果是业务校验不通过、消息本身有问题设true就会无限循环把MQ拖死这时候应该设false让消息进入死信队列后面人工或定时任务处理。RabbitListener默认会反序列化消息体如果你的消息体是JSON字符串需要配置Jackson2JsonMessageConverter否则Spring会用JDK默认序列化跨语言就完蛋了。3.4 JSON消息体为什么是首选这个话题很多人不重视直到踩了坑。如果两端都是Java服务用JDK序列化没毛病。但只要有一天你要从C#、Python、Go发送消息或者要消费别人发的消息JDK序列化直接废掉。所以我在项目里通常把消息体统一设计成JSON格式生产端发送前把对象转JSON字符串消费端收到字符串后再用Jackson反序列化。序列化配置Bean public MessageConverter messageConverter() { Jackson2JsonMessageConverter jsonMessageConverter new Jackson2JsonMessageConverter(); jsonMessageConverter.setCreateMessageIds(true); // 每条消息生成MessageId幂等和追踪都好用 return jsonMessageConverter; }配上这个Converter之后convertAndSend就会自动把对象转成JSON消费端也自动反序列化成目标类型。别忘了在application.yml里把spring.rabbitmq.listener.simple对应的converter配上或者直接在RabbitListenerContainerFactory里设置。4. 死信队列实战30分钟超时未支付是怎么被自动关掉的4.1 死信从哪里来死信队列并不神秘。三个来源消息被消费者拒绝basicNack且requeuefalse、消息TTL过期、队列达到最大长度。设计死信队列时先给业务主队列设置死信参数消息在过期或被拒绝之后转投到指定交换机再路由到死信队列。这个在订单场景里特别实用用户下单后创建一笔待支付订单30分钟内没支付就自动关单。实现思路不是起定时任务每秒钟扫一遍订单表而是让RabbitMQ的TTL机制帮你倒计时消息发到待支付队列设置x-message-ttl180000030分钟一到消息自动变成死信转到死信交换机由关单消费者接收处理。4.2 完整的死信链路配置Configuration public class RabbitDeadLetterConfig { // 业务交换机 Bean public DirectExchange orderExchange() { return new DirectExchange(order.exchange); } // 死信交换机 Bean public DirectExchange orderDeadExchange() { return new DirectExchange(order.dead.exchange); } // 待支付队列30分钟TTL死信转入死信交换机 Bean public Queue pendingOrderQueue() { MapString, Object args new HashMap(); args.put(x-message-ttl, 30 * 60 * 1000); args.put(x-dead-letter-exchange, order.dead.exchange); args.put(x-dead-letter-routing-key, order.timeout.dead); return new Queue(order.pending.queue, true, false, false, args); } // 死信队列 Bean public Queue timeoutOrderQueue() { return new Queue(order.timeout.dead.queue, true); } Bean public Binding pendingBinding() { return BindingBuilder.bind(pendingOrderQueue()).to(orderExchange()).with(order.created); } Bean public Binding deadBinding() { return BindingBuilder.bind(timeoutOrderQueue()).to(orderDeadExchange()).with(order.timeout.dead); } }消费者监听同一个死信队列order.timeout.dead.queue收到消息就去更新订单状态。这套方案比定时任务好在哪里订单量大的时候定时任务不知道哪些订单过期了每次都全表扫而RabbitMQ是把“哪些订单该关了”这个计算分布到消息过期机制里每个过期订单就是一个精准的触发信号。4.3 消费幂等消息可以被重复投递但订单不能重复关死信队列方案有一个副作用RabbitMQ投递消息不保证只投一次消费者在basicAck之前如果进程挂了消息会被重新投递。这意味着关单消息可能重复收到。所以必须在业务层面做幂等。我的做法是用Redis的setIfAbsent做一次性标记public boolean tryMarkOrderClosed(String orderId) { Boolean success redisTemplate.opsForValue() .setIfAbsent(order:closed: orderId, 1, Duration.ofDays(1)); return Boolean.TRUE.equals(success); }消费的时候先查这个标记如果已经处理过就直接basicAck跳过。数据库层面也可以给订单状态加一个“关单处理中”的唯一索引或状态机校验双保险。4.4 三十万个过期消息压测后的结论热搜词里有一条“rabbitmq 死信30分种会压多少”我实际压过一次。往待支付队列一次性塞30万条消息TTL都是30分钟。到点之后死信并不是瞬间全部转出来的RabbitMQ会尽可能均匀地把过期消息投递到死信队列但消费端如果只有一个消费者队列积压会非常明显吞吐量也就几百条每秒。结论是死信队列的消费者要开多实例并发而且每个实例的prefetch别设太大防止消息全堆积在消费者本地。如果你预估同时有过期订单的量级在几万以上建议关单消费者至少开5到10个并发实例。5. 跨语言调用C#生产者推送、Java消费者接收的坑点5.1 C#端怎么推送用.NET环境开发遇到最多的就是C#项目要往RabbitMQ推消息而消费端是Java的Spring Boot服务。C#这边用官方库RabbitMQ.ClientNuGet直接搜就能装。using RabbitMQ.Client; using System.Text; using Newtonsoft.Json; var factory new ConnectionFactory { HostName 127.0.0.1, Port 5672, UserName admin, Password admin123, VirtualHost /, AutomaticRecoveryEnabled true }; using var connection factory.CreateConnection(); using var channel connection.CreateModel(); var message new Dictionarystring, object { { orderId, 20250101000001 }, { amount, 99.9 } }; var props channel.CreateBasicProperties(); props.ContentType application/json; props.Persistent true; // 消息持久化 var body Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(message)); channel.BasicPublish( exchange: order.exchange, routingKey: order.created, basicProperties: props, body: body );AutomaticRecoveryEnabledtrue这个配置我强烈建议加上C#生产端如果和RabbitMQ断连它能自动重连不然后半夜服务重启一次就再也连不上了。5.2 跨语言最容易踩的坑属性名大小写我遇到过的最隐蔽的坑C#用Newtonsoft.Json序列化Dictionarystring, object时默认属性名是首字母大写OrderId、AmountJava端用Jackson反序列化时如果实体类字段是orderId、amountJackson默认区分大小写直接反序列化失败。解决办法有三个C#序列化时指定StringEscapeHandling或直接用JsonConvert.SerializeObject(message, new JsonSerializerSettings { ContractResolver new CamelCasePropertyNamesContractResolver() })输出小写属性名。Java端字段加JsonProperty(OrderId)注解但这种做法没法灵活应对字段名变化。更推荐C#定义一个DTO类而不是Dictionary直接控制属性名。还有一个坑是时间格式。C#的DateTime序列化出来是/Date(1700000000000)/这种格式Java的Jackson默认不认识。最好的办法是C#端统一序列化成字符串yyyy-MM-dd HH:mm:ss跨语言传时间永远用字符串。6. 集群与压测Docker三节点集群和一个关于死信压力的实测结论6.1 用Docker搭一个镜像模式集群单机模式跑通业务只是第一步生产环境至少需要三节点。Docker下搭集群的思路是先把三个RabbitMQ容器启动用同一个Erlang Cookie把它们组成集群再设置镜像队列策略让队列跨节点复制。# 启动第一个节点 docker run -d --name rabbitmq1 --hostname rabbitmq1 \ -e RABBITMQ_ERLANG_COOKIEcookiesecret \ -p 5672:5672 -p 15672:15672 \ rabbitmq:3-management # 启动第二个节点通过link访问第一个节点 docker run -d --name rabbitmq2 --hostname rabbitmq2 \ -e RABBITMQ_ERLANG_COOKIEcookiesecret \ --link rabbitmq1 \ rabbitmq:3-management # 在第二个节点容器内加入集群 docker exec -it rabbitmq2 rabbitmqctl stop_app docker exec -it rabbitmq2 rabbitmqctl join_cluster rabbitrabbitmq1 docker exec -it rabbitmq2 rabbitmqctl start_app第三个节点同理。注意RABBITMQ_ERLANG_COOKIE三个节点必须完全一样hostname不能重复否则加不进集群。节点组好之后还要设置镜像队列策略让队列在每个节点都有副本docker exec -it rabbitmq1 rabbitmqctl set_policy ha-all ^ {ha-mode:all,ha-sync-mode:automatic}这样任意一个节点宕机队列数据仍能从其他节点恢复。注意镜像队列有性能损耗对一致性要求极高的场景才用ha-modeall一般业务用ha-modeexactly配两个副本就够了。6.2 关于压测的几个实测结论压测别盲目凭感觉。我在测试环境压过10万条消息的积压场景几个结论分享给大家。单个消费者消费JSON消息不做复杂业务处理吞吐量大约在每秒3000到5000条。一旦加了数据库写入立刻掉到几百条。所以消费者线程里的业务逻辑越轻越好重活切出去异步处理。发布确认在高吞吐下也会成为瓶颈。publisher-confirm-typecorrelated每条消息都要等回调才能确认数据量大时吞吐下降明显。这时候可以把发送改成批量模式比如攒100条消息一次性basicPublish确认也等批量回调吞吐能提升好几倍。TTL死信积压的实测结论我在死信那一章说过了30分钟到期的消息会集中冲击死信队列消费端开多实例是最直接的解决方案。如果死信队列消费速度跟不上消息会攒在队列里占用内存最终触发RabbitMQ的流控机制整个节点都变慢。所以死信队列的消费者要预分配足够的并发量或者给死信队列也设置合理的TTL归档到别的存储系统。我自己的习惯是所有RabbitMQ核心队列都配合监控面板盯Ready和Unacked两个数字。Ready不断上涨说明消息在堆积Unacked高得离谱说明消费者处理不过来。死信队列也单独建一个面板这样出了问题不用翻日志凭感觉猜。写到最后还是那句话RabbitMQ的坑大多不在代码本身而在对它的消息模型的把握程度。先把交换机、队列、路由键这三件事想透再落实ACK、持久化、幂等这些细节剩下就是经验积累的问题了。本文还有配套的精品资源点击获取
返回列表