RabbitMQ五大消息模型详解:交换机、路由键与实战避坑指南
刚接触 RabbitMQ 的时候最劝退我的不是安装而是官方文档里那一堆概念交换机、绑定、路由键、虚拟主机……尤其是教程一口气列出的好几种消息模型我当年看一遍忘一遍看了三遍才敢说自己入门了。后来我换了种思路把几个模型按交换机类型重新理了一遍回头再看发现所谓的五大消息模型本质上就是一句话生产者不直接把消息扔进队列而是先交给交换机再由交换机按规则塞进一个或多个队列。这句话想通了五兄弟就认全了一半。这篇博文把 RabbitMQ 的五大消息模型讲透简单队列、工作队列、发布订阅、路由直连、主题通配符。每个模型都会带上原理、核心代码、注释说明以及我在真实业务里踩过的坑。如果你正在学 RabbitMQ或者项目里需要选型这篇文章可以直接当参考资料来用。1. 上菜之前先把 RabbitMQ 的“路由器逻辑”讲透很多新手第一次看 RabbitMQ 五个模型的时候看到一堆术语就懵了。我自己的经验是先别急着写代码先把那条消息的流动路径彻底搞清楚。RabbitMQ 里有一条铁律和传统消息队列很不一样——生产者永远不直接面对队列。消息的真实路径是这样的生产者 - 交换机(Exchange) - 绑定关系(Binding) - 队列(Queue) - 消费者你可以把交换机想象成快递分拣中心把队列想象成各个小区的快递柜。寄件人生产者把包裹消息丢给分拣中心分拣中心按面单路由键决定往哪个快递柜送。如果只有一个快递柜就只送到那一个如果要给所有小区都送一份就广播出去如果按“省内件”“同城件”分类就精确匹配路由。1.1 交换机的三种类型就是分拣中心的三种工作模式RabbitMQ 内置的交换机主要有四种类型其中三个是本文五个模型的核心交换机类型路由规则对应消息模型Fanout广播发给所有绑定的队列发布订阅模型Direct精确匹配路由键路由模型Topic通配符匹配路由键主题模型Headers按消息头匹配很少用一般不归入五大模型简单队列和工作队列其实用的是默认交换机default exchange这个默认交换机是一个空的 direct 交换机名字叫(AMQP default)它有个特性所有队列都会被它自动绑定到自己的队列名这个路由键上。所以我们在最简单的例子里直接把消息发到“名为 hello 的队列”本质上是发给默认交换机由默认交换机按队列名转投进去。这一点很多教程不讲但理解了它后面再学自定义交换机就顺了。1.2 五大模型一张表看透模型交换机队列数量核心场景简单队列默认交换机一收一发入门 Demo、单点任务工作队列默认交换机一个队列多个消费者任务分发、负载均衡发布订阅Fanout每个消费者一个队列广播通知路由Direct按路由键分流按级别分发日志主题Topic通配符匹配队列灵活分流、复杂业务事件这张表建议收藏。后面每一节都是围绕它展开的。2. 动手前必看Erlang 版本坑、安装步骤和启动失败自救手册先说环境因为根据我的经验RabbitMQ 劝退新手的第一关往往不是消息模型本身而是装好了启动不了。热搜里也有大量“rabbitmq启动失败”“rabbitmq erlang”这类词可见这是普遍痛点。RabbitMQ 是用 Erlang 写的所以它依赖 Erlang 运行环境。这个依赖关系是强绑定的不是随便装个最新版 Erlang 就能跑。不同版本的 RabbitMQ 要求不同版本的 Erlang这个兼容表在官网文档里专门有个页面叫 “RabbitMQ Erlang Version Requirements”你搜索这个关键词就能找到。2.1 下载安装的版本对应关系我实际用过的一个组合是RabbitMQ 3.12.x 要求 Erlang 26.0 以上RabbitMQ 3.13.x 要求 Erlang 26.2 以上。如果你在 Windows 上装建议直接用官方安装包它会自动带 Erlang 或者提示你装对应版本。Linux 上如果装了带有 Erlang 依赖的旧包很容易出现版本对不上。一个典型的启动失败现场是你启动了rabbitmq-server结果服务报错日志里出现类似Protocol inet_tcp not supported或者BOOT FAILED之类的字样。八成就是 RabbitMQ 和 Erlang 的版本组合不在兼容范围内。2.2 装完先干三件事装完且服务能启动之后我建议立刻做这三件事别急着写代码启动管理台插件rabbitmq-plugins enable rabbitmq_management访问管理台浏览器打开http://localhost:15672创建自己的账号默认 guest 账号只能通过 localhost 登录如果你是远程访问或者分布式环境guest 会直接拒绝你。创建一个管理员账号rabbitmqctl add_user myuser mypassword rabbitmqctl set_user_tags myuser administrator rabbitmqctl set_permissions -p / myuser .* .* .*这几条命令里set_permissions是给用户在虚拟主机/上配置权限的依次是配置权限、写权限、读权限.*表示全部权限。虚拟主机vhost这个概念以后用多了会理解现在先知道它是 RabbitMQ 里的“租户隔离”就够。2.3 启动失败排查三件套我整理了一个排查顺序按这个顺序来基本五分钟内能定位问题现象第一步看什么常见根因服务启动几秒后自动退出日志目录Windows 用户目录下AppData\Roaming\RabbitMQ\logLinux 在/var/log/rabbitmq/Erlang 版本不兼容、主机名解析失败5672 端口起不来netstat -ano | findstr 5672Windows或ss -lntp | grep 5672Linux端口被占用另一个 RabbitMQ 进程还活着管理台打不开先确认 15672 端口有监听再检查插件状态rabbitmq_management没启用主机名解析失败这一点很多人忽略。Linux 上如果你改过/etc/hostnameRabbitMQ 的 Erlang 节点名解析可能失败启动日志里会看到epmd相关的报错。解决方式很土但有效把本机主机名加进/etc/hosts映射到127.0.0.1。3. 模型一和模型二队列就是一根管子但竞争消费比想象中讲究两个模型放在一起讲因为它们用的是同一个东西默认交换机 队列。区别只在于消费者数量。3.1 简单队列一个消息只能被一个人吃掉最经典的 Hello World。生产者往队列发一条消息消费者从队列收一条消息。用 Java 原生客户端演示依赖就一个dependency groupIdcom.rabbitmq/groupId artifactIdamqp-client/artifactId version5.20.0/version /dependency生产者import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.Connection; import com.rabbitmq.client.Channel; public class Send { private final static String QUEUE_NAME hello; public static void main(String[] argv) throws Exception { // 1. 创建连接工厂设置 broker 地址 ConnectionFactory factory new ConnectionFactory(); factory.setHost(localhost); // 2. 建立 TCP 连接 try (Connection connection factory.newConnection(); Channel channel connection.createChannel()) { // 3. 声明队列queueDeclare(队列名, 是否持久化, 是否独占, 是否自动删除, 扩展参数) // 这里第二个参数 false表示队列不持久化重启后队列消失 channel.queueDeclare(QUEUE_NAME, false, false, false, null); String message Hello RabbitMQ!; // 4. 发送消息basicPublish(交换机名, 路由键, 附加属性, 消息体) // 这里交换机用空字符串表示使用默认交换机 channel.basicPublish(, QUEUE_NAME, null, message.getBytes()); System.out.println( [x] Sent message ); } // 5. 自动关闭连接 } }消费者import com.rabbitmq.client.*; public class Recv { private final static String QUEUE_NAME hello; public static void main(String[] argv) throws Exception { ConnectionFactory factory new ConnectionFactory(); factory.setHost(localhost); Connection connection factory.newConnection(); Channel channel connection.createChannel(); // 消费者这边也要声明队列幂等操作队列不存在才创建 channel.queueDeclare(QUEUE_NAME, false, false, false, null); System.out.println( [*] Waiting for messages. To exit press CTRLC); // 核心DefaultConsumer 是回调接口消息到达时触发 handleDelivery DeliverCallback deliverCallback (consumerTag, delivery) - { String message new String(delivery.getBody(), UTF-8); System.out.println( [x] Received message ); }; // autoAck 设为 true消费者收到消息后自动回执确认RabbitMQ 立即删除消息 channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag - { }); } }注意这里有个细节try-with-resources里 connection 和 channel 用完后会自动关闭。但这个写法只适合一次性发完消息就走人的场景。真实的生产者往往是常驻进程连接不能关这一点后面专门讲。3.2 工作队列能者多劳还是平均分配prefetch 是关键工作队列解决的问题是一个队列多个消费者怎么分消息我直接说结论默认情况下RabbitMQ 是把消息逐个轮流分发给每个消费者不管消费者处理得快还是慢。比如 10 条消息、两个消费者就是消费者 A 拿 1、3、5、7、9消费者 B 拿 2、4、6、8、10。这种叫轮询分发round-robin。问题来了。如果消费者 A 处理很慢消费者 B 处理很快A 那边还会积压。所以真实业务里几乎所有工作队列都要设置basicQos。// 消费者代码里加一行prefetch1 int prefetchCount 1; channel.basicQos(prefetchCount);basicQos(1)的含义是给这个消费者最多预取 1 条消息只有等这条消息 ack 之后RabbitMQ 才会再给它发下一条。这样处理快的消费者自然拿得多这就是“能者多劳”的公平分发机制。3.3 手动 ack 和持久化的三大配套设置上面简单队列里autoAck true是偷懒写法。生产上不建议开自动确认因为你没法控制消费者处理完消息的时机。如果消费者在自动 ack 之后、还没处理完就宕机了这条消息其实已经丢了。正确姿势是手动 ackDeliverCallback deliverCallback (consumerTag, delivery) - { try { String message new String(delivery.getBody(), UTF-8); System.out.println( [x] Received message ); // 模拟业务处理 doWork(message); } finally { // 手动确认参数要求 RabbitMQ 只确认当前这条消息 channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } }; channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag - { });配合手动 ack队列声明也要做持久化改造三个配套缺一不可配置项代码位置作用队列持久化queueDeclare(QUEUE_NAME, true, ...)队列本身在 broker 重启后不消失消息持久化basicPublish(..., MessageProperties.PERSISTENT_TEXT_PLAIN, ...)消息写入磁盘broker 重启不丢失手动确认autoAck falsebasicAck消费者挂掉后消息重新入队只有队列持久化、消息不持久化重启后队列在但消息全丢了。只有消息持久化、队列不持久化队列都没了消息自然也没了。这仨必须一起上才叫“至少一次的可靠投递”。注意PERSISTENT_TEXT_PLAIN 只是把消息的 deliveryMode 设为 2RabbitMQ 会按消息属性写入磁盘。但 RabbitMQ 的持久化并不是每条消息实时 fsync 的它有缓冲机制极端断电场景下还是可能丢数据这个要有心里预期。4. 模型三和模型四交换机一登场之前的写法全都要改前两个模型用的是默认交换机学的时候很方便但它掩盖了一个事实真正的路由能力在交换机。从发布订阅模型开始我们就要自己声明交换机、建立绑定关系了。4.1 发布订阅模型同一个消息广播给所有人想象一个场景用户下单成功后短信服务要发通知邮件服务要发邮件站内信服务要记一条记录。三个服务彼此独立但都得知道“下单成功”这个事件。这就是发布订阅的典型用法。用 Fanout 交换机实现广播核心逻辑是生产者只把消息发给交换机交换机扇出fan out给每一个绑定了它的队列。public class EmitLog { private static final String EXCHANGE_NAME logs; public static void main(String[] argv) throws Exception { ConnectionFactory factory new ConnectionFactory(); factory.setHost(localhost); try (Connection connection factory.newConnection(); Channel channel connection.createChannel()) { // 声明 exchangefanout 类型持久化设为 true channel.exchangeDeclare(EXCHANGE_NAME, fanout); String message 下单成功事件; // 注意路由键换成空字符串fanout 忽略路由键 channel.basicPublish(EXCHANGE_NAME, , null, message.getBytes()); System.out.println( [x] Sent message ); } } }消费者这边每个服务都要干三件事声明交换机、声明自己的专属队列、绑定。public class ReceiveLogs { private static final String EXCHANGE_NAME logs; public static void main(String[] argv) throws Exception { ConnectionFactory factory new ConnectionFactory(); factory.setHost(localhost); Connection connection factory.newConnection(); Channel channel connection.createChannel(); channel.exchangeDeclare(EXCHANGE_NAME, fanout); // queueDeclare() 不传队列名让 RabbitMQ 生成临时随机队列名 // exclusivetrue队列只对当前连接可见连接断开自动删除 String queueName channel.queueDeclare().getQueue(); // 绑定把临时队列绑到 logs 交换机上 channel.queueBind(queueName, EXCHANGE_NAME, ); System.out.println( [*] Waiting for messages. To exit press CTRLC); DeliverCallback deliverCallback (consumerTag, delivery) - { String message new String(delivery.getBody(), UTF-8); System.out.println( [x] Received message ); }; channel.basicConsume(queueName, true, deliverCallback, consumerTag - { }); } }这里的关键是channel.queueDeclare().getQueue()—— 不传任何参数RabbitMQ 会生成一个随机名字的临时队列比如amq.gen-JzTY20BRgKO-HjmUJj0wlg。这种队列的价值在于它天然是临时性的消费者挂了队列就没了不会堆积历史消息。对于广播场景这通常正是我们想要的——每个消费者只关心它在线期间的事件。4.2 路由模型按级别分发direct 交换机精确匹配发布订阅的问题是选择性太弱。假设我有一套日志系统错误日志要写数据库普通日志只要打印出来。Fanout 会把所有日志都塞给两个队列导致写数据库的服务收到大量无用消息。Direct 交换机解决这个问题的思路很直接队列绑定交换机时带上一个路由键生产者发消息时也带路由键两者精确相等才会路由过去。生产者public class EmitLogDirect { private static final String EXCHANGE_NAME direct_logs; public static void main(String[] argv) throws Exception { ConnectionFactory factory new ConnectionFactory(); factory.setHost(localhost); try (Connection connection factory.newConnection(); Channel channel connection.createChannel()) { channel.exchangeDeclare(EXCHANGE_NAME, direct); // 路由键可以是 info、warning、error 等 String severity error; String message 数据库连接池耗尽; channel.basicPublish(EXCHANGE_NAME, severity, null, message.getBytes()); System.out.println( [x] Sent severity : message ); } } }消费者绑定String queueName channel.queueDeclare().getQueue(); // 队列只接收 error 级别的消息 channel.queueBind(queueName, EXCHANGE_NAME, error);这里要注意的是交换机不存储消息队列也不存储被路由键过滤掉的消息。如果某个时刻没有队列绑定到error这个路由键这条消息就直接丢弃了。所以使用 direct 模型时保证队列和路由键先绑好再让生产者发消息。这也是 RabbitMQ 里一个非常经典的“消息丢了”的原因——你发得早绑定还没建立消息自然下落不明。4.3 关于“为什么不能直接把消息发到队列”的思考我见过不少同事刚接触 RabbitMQ 的时候尝试绕过交换机直接把消息 basicPublish 到某个队列名上。能成功因为默认交换机帮你转了。但一旦你希望“同一条消息按条件给不同的下游”直接发队列就做不到了。交换机的引入本质上是把“消息的目的地”和“生产者的发送行为”解耦生产者只关心自己发出的事件语义比如order.created至于这个事件最终触发多少服务、哪些服务关心它、什么时候有新服务加入——生产者一概不关心。新服务加入时只需要自己建队列、绑定到对应交换机生产端一行代码不用改。这就是发布订阅和路由模型真正的价值。4.4 Fanout 和 Direct 的代码 Altogether两步走1. 生产者声明交换机fanout/direct发送消息到交换机 2. 消费者建自己的队列把队列绑定到交换机fanout不写路由键direct写精确路由键如果是 Fanout多个消费者各绑各的队列消息复制多份如果是 Direct消费者只绑定自己感兴趣的精确路由键消息被过滤分发。5. 模型五主题模型路由键学会了通配符写作“#”和“*”主题模型Topic是 Direct 的升级版。它和 Direct 唯一的区别是路由键支持通配符匹配。Direct 只能精确匹配一个词Topic 可以模糊匹配。5.1 通配符规则两个符号记住就走不远RabbitMQ 的路由键是一个点分字符串比如order.created、user.login.failed。Topic 交换机匹配时有两个通配符可用通配符含义举例*匹配一个点分片段order.*匹配order.created不匹配order.created.success#匹配零个或多个片段order.#匹配order.created也匹配order.created.success注意*只能匹配“一个单词”而#可以匹配“多个单词”。比如路由键order.created.successorder.*不匹配order.#匹配。// 生产者发一条关于下单成功的主题消息 channel.basicPublish(topic_logs, order.created, null, message.getBytes());消费者 A只关心订单相关的所有消息。channel.queueBind(queueName, topic_logs, order.#);消费者 B只关心订单创建这个精确事件。channel.queueBind(queueName, topic_logs, order.created);消费者 C关心所有日志。channel.queueBind(queueName, topic_logs, #);#单独出现就是匹配所有路由键行为上接近 fanoutorder.*就是只匹配一层行为上接近 direct 但支持了一类前缀。所以业界常说 Topic 是集大成者这句话是有道理的。5.2 通配符匹配边界以及常见误解我第一次用 topic 的时候以为*能匹配任意的点分后缀结果调了半天不生效。其实*匹配“一个”点分单元不是“任意多个”。判断规则就一条把路由键按.拆分*占一个坑#占任意数量的坑。路由键绑定模式*.orange.*是否匹配绑定模式lazy.#是否匹配lazy.orange.rabbit是是quick.orange.fox是否lazy.brown.fox否第一个 是lazy但模式要求任意.此处*占brown后面还要一个 .而 fox 后面没有了是lazy.pink.rabbit.abc否超过三层是orange.rabbit否只有两段模式要三段否懒得记表的就记一句话*像正则里的[^.]#像正则里的.*但#必须出现在模式末尾或作为独立段出现RabbitMQ 允许a.#.b这种写法但官方不建议容易引起困惑我自己从不用。最常用的就是a.#、#.b、#这三种。5.3 实战中的路由键命名规范Topic 模式好不好用一半取决于路由键命名。我给自己定了一套规则也推荐给你用一面领域语言不要用技术语言。比如order.created、order.paid.success、user.login.failed看到名字就明白业务事件。层级控制在三段以内超过三段后通配符的意义就打折扣。不要以#开头做全局匹配路由键容易误伤。路由键里的语义要稳定一旦发布出去修改语义会导致历史消费者收到理解错误的消息。很多团队会把路由键设计成 “实体.动作.结果” 三段式这在电商场景里很好用。6. 五兄弟怎么选场景对照表以及“混用”才是常态看到了这里你可能会问那我项目里到底用哪个我的答案是别想着只用某一个模型。一个真实系统里五个模型往往同时存在只是分工不同。6.1 场景选型对照表业务诉求建议模型选择理由点对点任务比如生成报表工作队列多个 worker 分担压力一个事件触发多个独立服务发布订阅广播各管各的同一事件按类型分流处理路由精确匹配按级别隔离复杂业务事件多条件订阅主题一个交换机解决多种筛选请求/响应模式带返回值RPC 模型不属于五大模型但可以单独用6.2 一个订单创建能串起多少模型用一个典型电商场景展开用户下单成功。订单服务发出一个order.created主题消息发给order.topic交换机。物流服务用order.#绑定监听订单相关的所有变化包括创建、发货、签收。积分服务用order.created精确绑定只在创建时给用户加积分。通知服务用 fanout 交换机订单创建事件直接广播给短信、邮件、站内信三个队列。同时订单服务自己有个本地队列工作队列模式里面放着“创建订单快照”“推送数据仓库”这种任务交给多个 worker 慢慢消化。你看一个业务动作按不同语义走了三种交换机。这就是为什么我强烈不建议“只会一种模型走天下”——选型不是单选题是按消息的特征匹配分拣逻辑。6.3 选型前的三个灵魂拷问每次要引入新的消息通道时先问自己三个问题这条消息是否需要被多个独立下游处理是考虑 fanout否考虑 direct/work。相同的逻辑是否需要按类型过滤分发精确类型不懂模糊用 direct需要支持前缀/多级匹配用 topic。消费者是否必须消费到历史消息是队列必须持久化命名否可以直接用 RabbitMQ 生成的临时队列省心省力。这三个问题问完选型基本就定了。7. 上生产之前把这几颗雷先排了最后这部分是我在真实项目里挨过的最深刻的几颗雷代价是真金白银。写出来希望你不用再踩。7.1 重复声明队列参数不一致直接报错很多团队喜欢在消费者和生产者两边都queueDeclare同一个队列。如果参数一致没问题幂等。但一旦有一个人把durable从false改成true或加了某个参数RabbitMQ 会直接抛异常PRECONDITION_FAILED - inequivalent arg durable for queue xxx。原因是 RabbitMQ 里的队列声明是幂等但不是无脑幂等它要求参数完全一致。这个报错一旦出现最粗暴的办法是删掉原先队列再重新声明。但是——队列里如果有积压消息一删就全没了。所以我的建议是把队列声明统一收敛到一个配置类或者基础设施模块里全局只声明一次业务代码里不要动不动重复声明。7.2 Connection 和 Channel 的生命周期别用短连接新手最容易犯的错是把连接写在业务方法里每次发消息都 new 一个 Connection。Connection 是 TCP 长连接频繁创建销毁是一笔很大的开销而且在流量高的时候会成为瓶颈。正确做法是Connection 全局只建一次多个线程共享。每次发消息前可以从同一个连接里创建 Channel。Channel 是轻量级的会话通道创建成本低但也别每发一条就新建一个复用比频繁创建好得多。我用过一个稳定方案用一个连接池工具管理 Channel或者直接用 Spring Boot 的 RabbitTemplate它内部已经帮你做了连接复用。如果手写客户端至少做到连接单例、Channel 按条件复用。7.3 消费者忘记 ack消息积压是温水煮青蛙手动 ack 开了之后如果消费者的handleDelivery里抛了异常而且你只 catch 了业务异常没有处理 ack消息会一直处于 unacked 状态。表现是管理台上Unacked数字越来越大Ready的数字不动下游消费者明明在跑但就是觉得“消息走得慢”。这里要特别注意 try-finally 和处理失败的策略。如果是业务异常导致消息本身有问题那不能无限重试。一个从来都无法被消费掉的消息会反复进队列把后面所有正常消息堵死。我常用的手段是重试几次之后消费端手动 ack 掉这条消息同时把消息体转存到“死信队列”或者“异常表”里事后人肉排查。好在 RabbitMQ 本身支持死信队列用x-dead-letter-exchange参数声明队列可以把消费失败的消息自动踢到备胎队列这个后面有机会单独写一篇细说。7.4 最后再分享一个小技巧正规的 RabbitMQ 客户端里连接断开之后一般有自动重连机制但 Spring Boot 的RabbitTemplate在断连瞬间发消息会直接抛异常不会自动等重连。我的习惯是在发送代码外面包一层重试最多重试三次重试间隔指数退避一百毫秒、五百毫秒、两秒。public void sendWithRetry(String exchange, String routingKey, Object message) { int maxRetry 3; int delayMs 100; for (int i 1; i maxRetry; i) { try { rabbitTemplate.convertAndSend(exchange, routingKey, message); return; } catch (Exception e) { if (i maxRetry) { throw e; } try { Thread.sleep(delayMs); } catch (InterruptedException ex) { Thread.currentThread().interrupt(); throw new RuntimeException(发送被中断, ex); } delayMs * 5; } } }这个写法粗浅但能扛住 broker 短暂重启和网络抖动的场景。真实项目里我会再配合 sentinel 或者 resilience4j 做更精细的熔断不过思路都是一样的发送侧要有失败重试的预案。自己在本地实验的时候可以手动 kill 掉 RabbitMQ 进程再启动观察一下重试是否生效这是我最常用来验证非健壮性的手段。RabbitMQ 的五个消息模型说到底就是一套路由规则的递进默认交换机直连队列 - 扇形广播 - 直连精确匹配 - 主题通配符匹配。把这套链路想清楚看任何 RabbitMQ 的配置都不会再犯迷糊。代码注释里我写得比较啰嗦刻意保留了这些“为什么”希望对你也有帮助。