ARTICLE DETAIL

资讯详情

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

RabbitMQ三种交换机模式详解:fanout、direct、topic与生产实践

RabbitMQ三种交换机模式详解:fanout、direct、topic与生产实践 做消息中间件这一块RabbitMQ 工作模式是绕不开的硬知识点也是 Java 高级工程师面试里最常被追问的环节。上一篇文章我们把简单队列、Work Queues、ACK 确认和持久化讲透了这篇继续往下走重点拆解 fanout、direct、topic 三种交换机的工作模式顺带把生产环境里必踩的 Quorum Queue 和权限坑一起聊掉。这篇内容适合两类人一类是准备 Java 面试、正在背 RabbitMQ 八股文的同学另一类是已经在项目里引入 RabbitMQ但发现“能发能收”和“路由正确”完全两回事的开发者。全文不讲废话直接给原理、给代码、给排查路径你完全可以照着抄。1. 发布/订阅模式如何实现“一次发布处处广播”1.1 什么时候会用发布/订阅模式工作模式从简单到复杂的演进本质是对“消息该给谁消费”这个问题的回答力度不同。简单队列和 Work Queues 时代消息只有一条队列消费者之间是竞争关系谁先拿到谁消费。但业务里大量场景不是“抢任务”而是“广播事件”。最常见的例子就是支付成功后的多系统通知。订单系统支付成功产生一条消息积分系统要加积分、风控系统要做交易分析、短信服务要发通知。如果用 Work Queues三个消费者挂到同一队列上消息只会被其中一个抢到积分加了短信没发业务直接出事故。有人会说那我生产者分别往三个队列各发一次不就行了行但业务代码里就全是循环发送逻辑某一个目标队列发送失败还会把主流程拖垮。发布/订阅模式Publish/Subscribe就是为这种场景设计的。这里的核心角色是 fanout 交换机。交换机负责接收生产者消息并按规则分发到绑定好的队列fanout 的规则最简单粗暴——完全忽略 routingKey把消息复制给所有绑定到它身上的队列。一条消息进来广播给所有订阅方各消费方互不干扰谁消费失败也不影响别人。1.2 fanout 的匹配规则以及一个常见误解fanout 交换机的路由规则一句话就能讲完不看 routingKey不匹配 bindingKey消息无条件复制到所有绑定的队列。这个规则简单但恰恰是很多人踩坑的地方。我见过不少同事在 fanout 场景下还纠结 routingKey 怎么写甚至给每个队列绑定不同的 key指望它做过滤。结果折腾半天发现消息全都能收到因为 fanout 压根不参与匹配。要过滤消息应该用后面要讲的 direct 或 topic 交换机fanout 只有“广播”一个职责。打个比方fanout 交换机就是公司的全员邮件通讯组。发件人把邮件丢进通讯组组里每个成员都能收到一模一样的邮件发件人不需要知道成员具体是谁更不用逐个邮箱发一遍。1.3 代码实操消费者用临时队列生产者只管发布以纯 Java API 为例fanout 模式的核心步骤分三块。第一块消费者声明自己的专属队列。fanout 模式最适合用临时队列因为每个消费方只需要接收“订阅期间”的广播消息ConnectionFactory factory new ConnectionFactory(); factory.setHost(localhost); Connection connection factory.newConnection(); Channel channel connection.createChannel(); // 声明 fanout 交换机durable 设置为 true channel.exchangeDeclare(ex.pay.success, BuiltinExchangeType.FANOUT, true); // 不传参数声明临时队列队列名自动生成 String queueName channel.queueDeclare().getQueue(); // 绑定到交换机routingKey 传空字符串即可 channel.queueBind(queueName, ex.pay.success, );第二块生产者声明同一个交换机并发布消息。注意发布时 routingKey 也是空的channel.exchangeDeclare(ex.pay.success, BuiltinExchangeType.FANOUT, true); String message {\orderId\:\20250113001\,\status\:\PAID\}; channel.basicPublish(ex.pay.success, , null, message.getBytes(StandardCharsets.UTF_8));第三块消费者按正常方式从自己的临时队列里取消息消费。每个消费方各自独立声明队列和绑定就能各自收到一份完整的消息副本。这里要说清楚临时队列的语义queueDeclare()不带参数会生成一个非持久化、独占、自动删除的随机队列。消费者连接断开后队列自动删除非常适合 fanout 的“在线订阅”场景。但如果业务是订单通知这类“消费者不在线也要保住消息”的场景临时队列就是灾难必须改成手动声明固定名称的持久化队列。1.4 发布/订阅模式的经验笔记fanout 模式我用下来有一个特别深的体会多个消费组订阅同一个广播事件时一定要各自建队列不能共用一条队列。共用队列会退化成 Work Queues 的竞争消费模式两个组抢一条消息谁抢到谁消费另一个组永远收不到。另一个容易忽视的点是消费者的启动顺序。fanout 交换机本身不存储消息消息只会进入队列队列绑定的动作发生在消费者启动时。如果生产者先启动发消息消费者后启动再绑定中间发的消息就进了黑洞。所以在演示环境或自测脚本里先启消费者再启生产者避免“消息丢了”的假象。排查 fanout 是否生效最快的方法是打开 Management UI 的 Exchanges 页面点击交换机名进入详情直接看 Bindings 列表里挂了几条队列。如果队列数和你预期不符那问题一定出在绑定环节而不是发布环节。2. 路由模式用 routingKey 做精确筛选2.1 路由模式解决的痛点发布/订阅模式解决了“大家都收到”的问题但现实里“全收到”往往不等于“该收到”。一个典型的场景是日志系统error 级别的日志要进告警通道并且立刻通知值班人员info 和 warn 级别只要进归档存储就足够了。如果所有日志都用 fanout 广播告警消费者要自己过滤掉大量 info 日志带宽浪费不说每次有新级别日志接入消费者代码都要跟着改一遍维护成本非常高。direct 交换机路由模式就是干这个的它要求消息携带的 routingKey 与队列绑定的 bindingKey 完全一致消息才会投递到该队列。2.2 direct 的匹配规则以及最容易被忽略的细节direct 的匹配规则还有一个名字叫“完全匹配”按字面理解就行。消息的 routingKey 和队列的 bindingKey 必须逐字节相同才算命中。大小写不同不匹配长度不同不匹配多一个点号也不匹配。直接看这个表体会一下“精确”的含义发送的 routingKey队列绑定的是 order.error是否投递order.errororder.error是order.Errororder.error否大小写不同order.error.fatalorder.error否长度不一致errororder.error否单词数量不同direct 虽然追求精确但它本身是可以表现出广播效果的。一个队列可以绑定多个 key比如告警队列同时绑定 error 和 fatal那 error 和 fatal 消息都能进告警队列多个队列也可以绑定同一个 key比如归档队列和统计队列同时绑定 error那 error 消息会被复制到这两个队列。所以 direct 不是只能一对一路由它的“广播范围”完全由绑定关系画出来。2.3 代码实操日志分级处理的完整示例声明一个 direct 交换机然后让告警队列绑定 error 和 fatal归档队列绑定 info、warn、errorChannel channel connection.createChannel(); channel.exchangeDeclare(ex.log.direct, BuiltinExchangeType.DIRECT, true); // 告警队列绑定 error 和 fatal channel.queueDeclare(q.log.alert, true, false, false, null); channel.queueBind(q.log.alert, ex.log.direct, error); channel.queueBind(q.log.alert, ex.log.direct, fatal); // 归档队列绑定 info、warn、error channel.queueDeclare(q.log.archive, true, false, false, null); channel.queueBind(q.log.archive, ex.log.direct, info); channel.queueBind(q.log.archive, ex.log.direct, warn); channel.queueBind(q.log.archive, ex.log.direct, error);生产者发送时只需要指定日志级别对应的 routingKeychannel.basicPublish(ex.log.direct, error, MessageProperties.PERSISTENT_TEXT_PLAIN, 订单服务出现空指针异常.getBytes(StandardCharsets.UTF_8));这条消息会进入 q.log.alert 和 q.log.archive 两个队列因为 error 同时绑定在这两个队列上。而一条 warn 消息只会进归档队列不会打扰告警消费者。2.4 路由模式的实战细节与坑direct 模式最大的坑是“静默丢失”。绑定 key 和发送 key 不一致时RabbitMQ 不会报任何错误消息就像扔进了黑洞生产者那边回执照常成功消费者这边永远收不到。我排查过不少线上问题最后都定位到某个服务的 routingKey 拼写不一致。最稳妥的办法是在 Management UI 的 Exchange 详情页看一眼绑定列表再用一个测试生产者故意发送带 key 的消息确认投递路径。第二个坑是队列声明重复导致 406 PRECONDITION_FAILED。同一个 vhost 下一个队列被两处代码重复声明但参数不一致比如一处 durable 是 true另一处是 falsechannel 会直接报错关闭。解决办法是统一队列声明的参数配置一般建议把队列定义收拢到一个配置类里别在多个生产者和消费者的代码里各自声明。第三个坑是线上改造时的绑定变更。direct 模式下如果调整绑定关系老队列里的历史消息会继续按旧绑定消费完新消息才会走到新绑定上。上线前用 UI 里的 Queues 页签看下消息堆积数别刚切完路由就看到消费者把老消息按新逻辑处理导致业务异常。3. 主题模式用 * 和 # 做多维匹配3.1 主题模式解决的场景direct 能做到精确筛选但业务世界里消息类型往往是有层级关系的。以订单域为例可能有 order.created、order.paid、order.canceled还有 order.refund.applied、order.refund.rejected。如果每个事件类型都建一条绑定关系用 direct 会把绑定矩阵维护到崩溃。topic 交换机主题模式用一套通配符规则替代了大量精确绑定。它把 routingKey 看成由点号分隔的若干单词队列绑定 key 里可以用*和#表达模糊匹配*表示恰好匹配一个单词#表示匹配零个或多个单词举例说明队列绑定order.#可以收到 order、order.created、order.pay.success 等所有订单域事件队列绑定*.paid只能收到 order.paid 这类以 paid 结尾的单后缀事件队列绑定order.*能收到 order.created、order.paid但收不到 order.pay.success因为 pay.success 是两个单词。3.2 代码实操订单事件的多维订阅声明一个 topic 交换机队列 A 关心所有订单事件队列 B 只关心支付成功事件Channel channel connection.createChannel(); channel.exchangeDeclare(ex.order.topic, BuiltinExchangeType.TOPIC, true); channel.queueDeclare(q.order.all, true, false, false, null); channel.queueBind(q.order.all, ex.order.topic, order.#); channel.queueDeclare(q.order.paid, true, false, false, null); channel.queueBind(q.order.paid, ex.order.topic, *.paid); String event {\orderId\:\20250113002\,\amount\:199.00}; channel.basicPublish(ex.order.topic, order.pay.success, MessageProperties.PERSISTENT_TEXT_PLAIN, event.getBytes(StandardCharsets.UTF_8));这条 order.pay.success 消息q.order.all 能收到q.order.paid 收不到。如果发送的 routingKey 是 order.paid那么两个队列都能收到。一个交换机、两条绑定就实现了“订单全链路”和“支付结果”两套消费视角这是 direct 模式很难优雅做到的。3.3 主题模式三大易错点第一路由键必须是单词加点号分隔的结构。RabbitMQ 对 routingKey 的约束包括最长 255 字节不能以点号开头或结尾。如果你在业务里直接把订单号当 routingKey一旦订单号超过 255 字节或者含特殊字符交换机匹配可能异常。第二*只匹配一个单词不是“任意一段”。order.* 匹配不了 order.pay.success这是主题模式里最高频的误解。很多人以为 * 是“通配任意长度”实际上它是“通配恰好一个单词”。第三#可以匹配零个单词。队列绑定 order.#不仅能收到 order.created也能收到一个 routingKey 就是 order 的消息甚至绑定#的队列能收到所有发给该交换机的消息相当于 fanout 的效果。这个“零个单词”的语义经常被忽略绑定时心里要有数。主题模式下还有个性能考量绑定关系特别多的时候topic 的匹配计算比 direct 重一些。常规业务几百条绑定完全无感但如果搞到上万条绑定比如做 IoT 设备级路由就要关注匹配耗时。这种量级的场景我会选择在应用层做路由而不是压在 topic 交换机上。4. 三种模式对比与选型经验4.1 一张表看懂三种交换机把前面三章的内容压缩成一张对比表面试和实际选型都能直接用交换机类型路由规则典型场景误用示例fanout忽略 routingKey广播到所有队列支付成功通知多系统、全员广播想做过滤却选 fanoutdirectroutingKey 与 bindingKey 完全匹配日志分级、按错误码分发想按事件类型层级匹配却选 directtopic用 * 和 # 模式匹配单词订单域事件分发、多维路由路由键不规范单词划分混乱4.2 选型心法先回答三个问题我在项目里给团队定过一个最简单的选型流程写代码前先问三个问题。第一消息需不需要被筛选完全不需要所有订阅方都要收到直接选 fanout。第二筛选条件是精确值还是模糊规则精确的、可枚举的路由键比如 error、success、refund选 direct 最简单直观。第三路由键是否有层级结构或者需要按模式归类比如 order.created、order.pay.success 这种多级事件选 topic 最合适。这里要给一个反向提醒别为了“灵活”盲目选 topic。topic 的能力最强但它对 routingKey 的命名规范要求也最高。团队如果没人维护路由键的命名约定时间一长绑定关系里全是 order.#、#.paid 这种拍脑袋的写法排查成本远高于 direct。在小团队里direct 反而是稳定性最好的选择因为规则简单不容易被写坏。4.3 顺势聊下 Kafka、RocketMQ 的选型视角既然是高级工程师视角提到队列选型就绕不开 RabbitMQ、Kafka、RocketMQ 的横向对比。我的经验是这不是“哪个更好用”的问题而是“你的消息需要什么语义”的问题。RabbitMQ 的优势在于复杂的路由能力也就是这篇文章讲的三种工作模式它更适合业务系统内部的事件解耦比如订单状态变更、通知分发开发效率很高消息可靠性机制也足够细腻。Kafka 的强项是超高吞吐和日志回放适合埋点日志、用户行为流、大数据链路但它的路由能力很弱基本就是按 topic 分发做不了 topic 模式那种细粒度模糊匹配。RocketMQ 在事务消息、定时消息、消息重试上做得最完善适合交易链路中对消息可靠性和时序要求极高的场景。所以面试如果被问到“为什么不用 Kafka 做业务通知”回答思路不是“Kafka 不好”而是“Kafka 没有 RabbitMQ 这种灵活的交换机路由模型而且业务事件的消费延迟在毫秒级RabbitMQ 更合适”。选型先看语义再看性能最后看运维成本。5. 生产环境进阶Quorum Queue 与可靠投递5.1 为什么镜像队列不再是首选工作模式讲清楚了生产环境还有一个绕不开的话题——队列本身的高可用和可靠性。以前做 RabbitMQ 高可用大家用的是镜像队列Mirrored Queue主节点接收消息后全量复制到从节点。镜像队列的问题在于主从全量复制开销大扩容不灵活而且脑裂风险高极端情况下可能出现丢消息。RabbitMQ 3.8 开始主推 Quorum Queue基于 Raft 协议实现多数派写替代镜像队列成了新项目的默认选择。你看热词里频繁出现 “quorum queue”就知道这个话题在面试里有多高频。5.2 如何创建一个 Quorum Queue创建 Quorum Queue 只需要在声明队列时加一个参数指定 x-queue-type 为 quorumMapString, Object args new HashMap(); args.put(x-queue-type, quorum); Channel channel connection.createChannel(); channel.queueDeclare(q.order.pay, true, false, false, args);队列创建后绑定工作模式的路由关系完全不受影响。你可以把 q.order.pay 绑定到 fanout、direct 或者 topic 交换机上Quorum Queue 对这三种工作模式一视同仁区别只在队列内部的复制和存储机制。用 docker 部署 RabbitMQ 后可以在命令行执行rabbitmqctl list_queues name type messages确认队列类型已经是 quorum。5.3 Quorum Queue 的限制和使用注意点Quorum Queue 带来了可靠性也带来了一些限制使用前要清楚。它不支持独占队列、不支持自动删除队列、不支持事务。消息只能是持久化的也就是投递模式必须为 2非持久化消息在 Quorum Queue 里其实是无效的。消费方式强烈建议手动 ACK因为 Quorum Queue 的投递语义和确认机制配合手动 ACK 才是最稳的。从镜像队列迁移到 Quorum Queue建议按这个步骤走新声明 Quorum Queue 队列绑定原来的路由关系切换消费者到新队列观察消息消费正常后再删除旧队列。别直接删旧队列一旦回滚路径没了线上出问题很难收场。5.4 可靠投递的组合拳publisher confirm 手动 ACK工作模式、Quorum Queue 都只是载体真正保证消息不丢的是一套组合配置。消息从产生到消费中间有四个可能丢失的环节发送到交换机、交换机投递到队列、队列存储、消费者处理。发送到交换机这一环用 publisher confirm 机制保证。在 channel 上调用confirmSelect()然后发送消息RabbitMQ 会返回确认回执异步监听handleAck和handleNack就能知道哪些消息投递失败了。交换机投递到队列这一环需要把消息的 mandatory 参数设为 true并注册 ReturnListener当没有队列接收消息时回调返回。队列存储这一环靠交换机、队列、消息三层都设置持久化。消费者处理这一环用手动 ACK 保证处理成功后再确认。这四步组合起来才是生产环境里真正“端到端不丢”的完整链路。单独开 Quorum Queue或者单独开 publisher confirm都只是解决了一个环节别指望某一步能包打天下。6. 生产环境常见问题速查与排查实录6.1 docker 部署后admin 账号无法创建虚拟主机这个坑在热词里反复出现我在项目里也真实遇到过。docker 部署 RabbitMQ 时设置环境变量 RABBITMQ_DEFAULT_USERadmin、RABBITMQ_DEFAULT_PASSadmin123能正常登录 Management UI但想通过 UI 创建新的 virtual host 时界面提示没有权限。原因很简单创建虚拟主机需要用户具备 administrator 标签而默认创建的 admin 用户可能只被标记为 management 标签只能管理自己的队列和交换器没有 vhost 级管理权限。解决办法是用命令行工具提升权限# 进入 docker 容器 docker exec -it rabbitmq bash # 给 admin 用户设置 administrator 标签 rabbitmqctl set_user_tags admin administrator # 确认用户标签 rabbitmqctl list_users如果你需要新建专门的业务虚拟主机还要用 set_permissions 给用户授权rabbitmqctl add_vhost /order_service rabbitmqctl set_permissions -p /order_service admin .* .* .*这里想说个经验权限配置别图省事全给.*。生产环境每个服务面对自己的 vhost只应该分配自己队列和交换机通配符的权限。比如积分服务只能对 /point_service 下以point.开头的队列有读写权限。一开始嫌麻烦后面出了乱子更麻烦。6.2 消息发布成功消费者就是收不到这个问题排在 RabbitMQ 求助榜前列。排查路径我按顺序说先确认生产者有没有开 publisher confirm如果开了确认回执是不是 ack这一步排除“消息压根没到 RabbitMQ”然后看 Management UI 的 Queues 页面找到目标队列看 Message rates 有没有入队数量如果入队数为零说明消息没有路由到这条队列最后点进交换机详情查看 Bindings 列表通常问题就出在 bindingKey 和发送的 routingKey 不一致。这里分享一个我自己的笨办法。排查绑定关系时不要直接用 Java 代码猜在 Management UI 里手动创建一个临时测试队列并绑定一个测试 key用发布者发一条测试消息看能不能进入测试队列能进去就说明交换机到队列的链路是通的问题在业务代码里的 key 写错了。6.3 队列声明报 406 PRECONDITION_FAILED这个报错的含义是队列已经存在但声明的参数和之前不一致。通常发生在同一个 vhost 下一个队列被多处代码声明一处设了 durabletrue另一处设了 durablefalse或者一处设置了 x-queue-typequorum另一处没设置。RabbitMQ 只允许队列“首次声明”时定义参数后续声明必须完全一致。最快的解决办法是删除旧队列重新声明但要注意清空队列里的积压消息。如果是生产环境建议保留旧队列换一个新队列名接入新配置切换完成后确认没有消费端依赖旧队列后再删。这也是为什么我前面强调队列配置要收拢到一个类里管理散落在代码里迟早出这种问题。6.4 虚拟主机和权限导致消费者连接失败消费者启动时报 vhost 访问拒绝常见原因是用户没有对应 vhost 的权限或者用户的 tag 只有 management 但不具备消息收发所需的 configuration、write、read 权限。排查命令是rabbitmqctl list_permissions -p /order_service正常输出应该看到用户对该 vhost 有配置、写、读三项权限。如果有权限但还是连不上检查连接配置里的 vhost 是否写错了默认是/如果业务用的是/order_service连接工厂里要把factory.setVirtualHost(/order_service)加上。这个问题在 Spring Boot 配置里特别常见因为 yml 里漏写 virtual-host 时默认连的永远是/业务 vhost 里的队列自然一个都消费不到。最后分享一点个人的体会这篇文章写到这里其实我最想说的是RabbitMQ 的三种工作模式本身并不复杂难的是绑定关系的管理和团队规范的建立。我在项目里吃过不少亏所以现在无论项目大小我都要先画一张消息拓扑图把交换机、队列、绑定 key、消费方全部标清楚再让代码按图实现。交换机命名统一加业务前缀比如 ex.order.pay、ex.log.direct队列也按用途命名q.order.paid、q.log.alert这样排查问题的时候光看名字就知道这条消息的完整路径。下篇我会沿着这个方向继续扩展把死信队列、延迟队列、消费端幂等这几个高频实战点逐个拆开讲这些都是工作模式落地时必然会遇到的进阶问题。
返回列表