官方文档:RabbitMQ Tutorials | RabbitMQ
Ⅰ. 7种工作模式介绍
一、Simple(简单模式)

-
特点: 一个生产者P,一个消费者C,消息只能被消费一次,也称为点对点(Point-to-Point)模式。
-
适用场景: 消息只能被单个消费者处理
二、Work Queue(工作队列模式)

一个生产者P,多个消费者C1、C2。在多个消息的情况下,Work Queue会将消息分派给不同的消费者,每个消费者都会接收到不同的消息 。RabbitMQ 默认使用 轮询 (Round-Robin)分发给消费者。
-
特点: 消息不会重复,分配给不同的消费者。
-
适用场景: 集群环境中做异步处理
- 比如 12306 短信通知服务,订票成功后,订单消息会发送到 RabbitMQ,短信服务从 RabbitMQ 中获取订单信息,并发送通知信息(在短信服务之间进行任务分配)

三、Publish/Subscribe(发布/订阅模式)

一个生产者P,多个消费者C1、C2,此外 X 代表交换机,它可以将消息复制多份,让多个消费者接收相同的消息 。
生产者发送一条消息,经过交换机转发到多个不同的队列,多个不同的队列就有多个不同的消费者。
适合场景 :消息需要被多个消费者同时接收的场景。比如:实时通知或者广播消息。
比如中国气象局发布 "天气预报" 的消息送入交换机之后,新浪、百度、搜狐、网易等门户网站接入消息,通过队列绑定到该交换机,自动获取气象局推送的气象数据。
概念介绍🐔
Exchange :交换机(X)
作用: 生产者将消息发送到 Exchange,由交换机将消息按一定规则路由到一个或多个队列中(上图中生产者将消息直接投递到队列中,实际上这个在 RabbitMQ 中不会发生)
Exchange 只负责转发消息,不具备存储消息的能力 ,因此如果没有任何队列与 Exchange 绑定,或者没有符合路由规则的队列,那么消息就会丢失。
RabbitMQ 交换机有四种类型:fanout、direct、topic、headers,不同类型有着不同的路由策略。
-
Fanout :广播策略 ,将消息交给所有绑定到该交换机的队列(Publish/Subscribe 发布订阅模式 )
-
Direct :定向策略 ,把消息交给符合指定
routing key的队列(Routing 路由模式 ) -
Topics :通配符策略 ,把消息交给符合
routing pattern的队列(Topics 通配符模式 ) -
headers :该类型的交换器不依赖于路由键的匹配规则来路由消息,而是根据发送的消息内容中的
headers属性进行匹配。(headers类型的交换器性能会很差,而且也不实用,基本上不会看到它的存在)
Routing Key :路由键。生产者将消息发给交换器时,指定的一个字符串,告诉交换机应该如何处理这个消息(即告诉交换机将该消息发送到哪里去)
Binding Key :绑定键。交换器与队列通过 Binding Key 关联起来。队列在绑定交换机的时候一般会指定一个 Binding Key,告诉交换机要接收哪些消息。

比如下图:如果在发送消息时,设置了
Routing Key为 orange,那么消息就会路由到 Q1。
当消息的
Routing key与队列绑定的Binding key相匹配时,消息才会被路由到这个队列 。此外,
Binding Key其实也属于路由键中的一种,官方解释为: the routing key to use for the binding。可以翻译为:在绑定的时候使用的路由键。大多数时候,包括官方文档和 RabbitMQ Java API 中都把
Binding Key和Routing Key看作Routing Key,为了避免混淆,可以这么理解:
生产者在发送消息的时候,需要的路由键是
Routing Key队列在绑定交换机的时候,需要的路由键是
Binding Key
四、Routing(路由模式)

路由模式是发布订阅模式的变种,在发布订阅模式的基础上,增加了 Routing Key。
发布订阅模式是无条件的将所有消息分发给所有消费者,路由模式是 Exchange 根据 Routing Key 的规则,将数据筛选后发给对应的消费者队列。
适合场景 :需要根据特定规则分发消息的场景。
比如系统打印日志,日志等级分为 error、warning、info、debug,就可以通过这种模式,把不同的日志发送到不同的队列,最终输出到不同的文件。
五、Topics(通配符模式)

路由模式的升级版,在 Routing Key的基础上,增加了通配符的功能 ,使之更加灵活。
Topics 和 Routing 模式的基本原理相同,即:生产者将消息发给交换机,交换机根据 Routing Key 将消息转发给与 Routing Key 匹配的队列。
不同之处:
-
Routing 模式是 精确匹配
-
Topics 模式是 模糊匹配 (通配符匹配)
适合场景 :需要灵活匹配和过滤消息的场景。
六、RPC(RPC通信模式)

在 RPC 通信的过程中,没有生产者和消费者,比较像 RPC 远程过程调用,就是通过两个队列实现了一个可回调的过程。
该模式通常用于 "客户端发送请求给服务端执行某个任务,并等待结果返回 "。换句话说,用异步消息机制模拟 "同步调用" 。
客户端发送消息到一个指定的队列,并在消息属性中设置
replyTo字段,这个字段指定了一个回调队列,用于接收服务端的响应。服务端接收到请求后,处理请求并发送响应消息到
replyTo指定的回调队列。客户端在回调队列上等待响应消息。一旦收到响应,客户端会检查消息的
correlationId属性,以确保它是所期望的响应。
适用场景: 在分布式系统中,我们常常有这种需求:"我有一个计算密集型任务(比如生成报告、图片识别、机器学习推理),我不想在主系统中直接做,而想交给后端的工作进程去做,然后拿到结果。"
七、Publisher Confirms(发布确认机制)

Publisher Confirms 模式是 RabbitMQ 提供的一种确保消息可靠发送到 RabbitMQ 服务器的机制。在这种模式下,生产者可以等待 RabbitMQ 服务器的确认,以确保消息已经被服务器接收并处理 ,从而避免消息丢失的问题。
-
生产者将
Channel设置为confirm模式后(通过调用channel.confirmSelect()完成),发布的每一条消息都会获得一个唯一ID,生产者可以将这些序列号与消息关联起来,以便跟踪消息的状态。 -
当消息被 RabbitMQ 服务器接收并处理后,服务器会异步地 向生产者发送一个确认
ACK给生产者(包含消息的唯一ID),表明消息已经送达。
适用场景 :对数据安全性要求较高的场景。比如金融交易、订单处理等等。
💥注意事项:
Publisher Confirms(发布确认机制)属于可靠层,与发布订阅、路由、主题等模式不冲突。
它们是 "并行概念",一个负责消息投递的可靠性(安全性) ,另一个负责消息分发的逻辑(路由模式) 。两者是互补 而不是互斥的。
所以完全可以在发布订阅(Fanout)等模式下,开启
Confirm确认机制来确保消息可靠投递 。
Ⅱ. 工作模式的使用案例
一、Work Queue(工作队列模式)

简单模式的增强版,和简单模式的区别就是:简单模式只有一个消费者,而工作队列模式支持多个消费者接收消息,消费者之间是竞争关系,每个消息只能被一个消费者接收。
首先先把之前代码简化一下,把常量抽出来:
public class Constants {
public static final String *IP *= "127.0.0.1";
public static final int *PORT *= 5672;
public static final String *VIRTUALHOST *= "lirendada";
public static final String *USERNAME *= "liren";
public static final String *PASSWORD *= "123456";
public static final String *WORK_QUEUE *= "work.queue";
}为了能看到多个消费者竞争的关系,这里一次发送10条消息。生产者代码如下所示:
public class producer {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂*
* *ConnectionFactory factory = new ConnectionFactory();
* *factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
*// 2. 创建连接Connection和通道*
* *Connection connection = factory.newConnection();
* *Channel channel = connection.createChannel();
*// 3. 声明一个队列*
* *channel.queueDeclare(**Constants.** ***WORK_QUEUE** *, true, false, false, null);
*// 4. 发送消息(当使用内置交换机的时候,routingKey必须和队列名称保持一致)*
* *for(int i = 0; i < 10; ++i) {
String text = "hello workqueue " + i;
channel.basicPublish("", **Constants.** ***WORK_QUEUE** *, null, text.getBytes(StandardCharsets.*UTF_8*));
}
*// 5. 释放资源*
* *connection.close();
}
}消费者代码和简单模式一样,只是复制两份,两个消费者代码可以是一样的:
public class consumer1 {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂*
* *ConnectionFactory factory = new ConnectionFactory();
* *factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
*// 2. 创建连接Connection和通道*
* *Connection connection = factory.newConnection();
* *Channel channel = connection.createChannel();
*// 3. 声明一个队列(这是安全性措施,因为如果生产者还没创建队列的话,消费者这边直接读取会报错)*
* *channel.queueDeclare(**Constants.** ***WORK_QUEUE** *, true, false, false, null);
*// 4. 接收消息,进行消费💥*
* *DefaultConsumer defaultConsumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body) throws IOException
{
System.*out*.println("接收到消息: " + new String(body));
}
};
channel.basicConsume(**Constants.** ***WORK_QUEUE** *, true, defaultConsumer);
}
}
二、Publish/Subscribe(发布/订阅模式)

RabbitMQ 交换机常见三种类型:fanout、direct、topic,不同类型有着不同的路由策略。
-
Fanout :广播策略 ,将消息交给所有绑定到该交换机的队列(Publish/Subscribe 模式 )
-
Direct :定向策略 ,把消息交给符合指定
routing key的队列(Routing 模式 ) -
Topic :通配符策略 ,把消息交给符合
routing pattern的队列(Topics 模式 )
所以发布订阅模式使用了交换机的 Fanout 广播模式来完成!此时需要知道创建交换机,以及绑定交换机和队列的方法,如下所示:
-
创建交换机:
channel.**exchangeDeclare** (String exchange, BuiltinExchangeType type, boolean durable, boolean autoDelete, boolean internal, Map<String, Object> arguments) throws IOException;参数 说明 典型取值/建议 exchange 交换机名称 如 "my.direct.exchange"。同名交换机若存在,属性必须一致type 交换机类型 DIRECT,FANOUT,TOPIC,HEADERSdurable 是否持久化 true表示 RabbitMQ 重启后交换机仍保留autoDelete 是否自动删除 当没有队列绑定 且没有连接使用 时自动删除 internal 是否内部交换机 若为 true,客户端不能直接发布消息 到该交换机,通常用于交换机间路由arguments 扩展参数 例如备用交换机、延迟特性等 -
绑定交换机和队列:
channel.**queueBind** (String queue, String exchange, String routingKey) throws IOException;参数名 作用 说明 queue 队列名称 要绑定的队列名(必须已经声明过 queueDeclare()) exchange 交换机名称 要绑定到的交换机(必须已声明过 exchangeDeclare()) routingKey 路由键 决定消息如何路由到该队列
下面是需要用到的常量:
*// 发布订阅模式*
public static final String *FANOUT_QUEUE1 *= "fanout.queue1";
public static final String *FANOUT_QUEUE2 *= "fanout.queue2";
public static final String *FANOUT_EXCHANGE *= "fanout.exchange";因为广播模式就是要推送消息给所有绑定当前交换机的队列,所以绑定队列和交换机的时候,只需要设置 routing key为空字符串即可 。
生产者代码:
public class producer {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂*
* *ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
*// 2. 创建连接Connection和通道*
* *Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
***// 3. 声明两个队列** *
* *channel.queueDeclare(Constants.*FANOUT_QUEUE1*, true, false, false, null);
channel.queueDeclare(Constants.*FANOUT_QUEUE2*, true, false, false, null);
***// 4. 创建交换机** *
* *channel.exchangeDeclare(Constants.*FANOUT_EXCHANGE*, BuiltinExchangeType.*FANOUT*, true, false, false, null);
***// 5. 绑定交换机和队列** **(因为是广播模式,所以本质不需要routingkey,置为空字符串即可)*
* *channel.queueBind(Constants.*FANOUT_QUEUE1*, Constants.*FANOUT_EXCHANGE*, "");
channel.queueBind(Constants.*FANOUT_QUEUE2*, Constants.*FANOUT_EXCHANGE*, "");
***// 6. 发送消息** *
* *String text = "hello public/subscribe && fanout mode!";
channel.basicPublish(Constants.*FANOUT_EXCHANGE*, "", null, text.getBytes(StandardCharsets.*UTF_8*));
*// 7. 释放资源*
* *connection.close();
}
}消费者代码:(有两份,都是一样的,这里只放一份)
public class consumer1 {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂*
* *ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
*// 2. 创建连接Connection和通道*
* *Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
***// 3. 声明队列** **(这是安全性措施,因为如果生产者还没创建队列的话,消费者这边直接读取会报错)*
* *channel.queueDeclare(Constants.*FANOUT_QUEUE1*, true, false, false, null);
*// 4. 接收消息,进行消费💥*
* *DefaultConsumer defaultConsumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body) throws IOException
{
System.*out*.println("接收到消息: " + new String(body));
}
};
channel.basicConsume(**Constants.** ***FANOUT_QUEUE1** *, true, defaultConsumer);
}
}三、Routing(路由模式)

RabbitMQ 交换机常见三种类型:fanout、direct、topic,不同类型有着不同的路由策略。
-
Fanout :广播策略 ,将消息交给所有绑定到该交换机的队列(Publish/Subscribe 模式 )
-
Direct :定向策略 ,把消息交给符合指定
routing key的队列(Routing 模式 ) -
Topic :通配符策略 ,把消息交给符合
routing pattern的队列(Topics 模式 )
路由模式采用的是 RabbitMQ 中的 Direct 定向策略,生产者发送消息的时候,交换机需要根据消息中的 Routing Key将消息发送给指定的队列 ,而不是发给每一个队列了!
此时,队列和交换机的绑定,不能是任意的绑定了,而是要指定一个 Binding Key。
只有队列绑定时的 Binding Key和消息中的 Routing Key完全一致,队列才会接收到消息 。

和发布订阅模式的区别是:交换机类型不同、绑定队列的 Binding Key 不同。
下面是需要用到的常量:
*// 路由模式
public static final String DIRECT_EXCHANGE = "direct.exchange";
public static final String DIRECT_QUEUE1 = "direct.queue1";
**public static final String DIRECT_QUEUE2 = "direct.queue2";*生产者代码:
public class producer {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂*
* *ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
*// 2. 创建连接Connection和通道*
* *Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
*// 3. 声明两个队列*
* *channel.queueDeclare(**Constants.** ***DIRECT_QUEUE1** *, true, false, false, null);
channel.queueDeclare(**Constants.** ***DIRECT_QUEUE2** *, true, false, false, null);
*// 4. 声明交换机*
* *channel.exchangeDeclare(**Constants.** ***DIRECT_EXCHANGE** *, BuiltinExchangeType.***DIRECT** *, true, false, false, null);
***// 5. 绑定交换机和队列** *
* *channel.queueBind(Constants.*DIRECT_QUEUE1*, Constants.*DIRECT_EXCHANGE*, "orange");
channel.queueBind(Constants.*DIRECT_QUEUE2*, Constants.*DIRECT_EXCHANGE*, "black");
channel.queueBind(Constants.*DIRECT_QUEUE2*, Constants.*DIRECT_EXCHANGE*, "green");
*// 6. 发送消息*
* *String text1 = "hello routing, i am orange!";
channel.basicPublish(Constants.*DIRECT_EXCHANGE*, **"orange"** , null, text1.getBytes(StandardCharsets.*UTF_8*));
String text2 = "hello routing, i am black!";
channel.basicPublish(Constants.*DIRECT_EXCHANGE*, **"black"** , null, text2.getBytes(StandardCharsets.*UTF_8*));
String text3 = "hello routing, i am green!";
channel.basicPublish(Constants.*DIRECT_EXCHANGE*, **"green"** , null, text3.getBytes(StandardCharsets.*UTF_8*));
*// 7. 释放资源*
* *connection.close();
}
}
消费者代码:(有两份,除了绑定队列不同外,基本都是一样的,这里只放一份)
public class consumer1 {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂*
* *ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
*// 2. 创建连接Connection和通道*
* *Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
*// 3. 声明队列(这是安全性措施,因为如果生产者还没创建队列的话,消费者这边直接读取会报错)*
* *channel.queueDeclare(Constants.*DIRECT_QUEUE1*, true, false, false, null);
*// 4. 接收消息,进行消费💥*
* *DefaultConsumer defaultConsumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body) throws IOException
{
System.*out*.println("接收到消息: " + new String(body));
}
};
channel.basicConsume(Constants.*DIRECT_QUEUE1*, true, defaultConsumer);
}
}
四、Topics(通配符模式)

Topics 和 Routing 模式的区别是:
-
交换机类型不同: Topics 模式使用的交换机类型为
topic;Routing 模式使用的交换机类型为direct。 -
匹配规则不同:
topic类型的交换机在匹配规则上进行了扩展,Binding Key支持通配符匹配;direct类型的交换机路由规则是Binding Key和Routing Key完全匹配。
匹配规则有如下要求:
-
Routing Key是一系列由点.分隔的单词 ,比如 "stock.usd.nyse"、"nyse.vmw"、"quick.orange.rabbit" -
Binding Key和Routing Key一样,也是点.分割的字符串 -
Binding Key中可以存在两种特殊字符串,用于模糊匹配-
*:表示一个单词 -
#:表示多个单词(0-N个)
-
比如:
-
Binding Key 为 "d.a.b" 会同时路由到 Q1 和 Q2
-
Binding Key 为 "d.a.f" 会路由到 Q1
-
Binding Key 为 "c.e.f" 会路由到 Q2
-
Binding Key 为 "d.b.f" 会被丢弃,或者返回给生产者(需要设置 mandatory 参数)
下面是需要用到的常量:
*// 通配符模式
public static final String TOPIC_EXCHANGE = "topic.exchange";
public static final String TOPIC_QUEUE1 = "topic.queue1";
**public static final String TOPIC_QUEUE2 = "topic.queue2";*生产者代码如下所示:
public class producer {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂*
* *ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
*// 2. 创建连接Connection和通道*
* *Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
*// 3. 声明队列*
* *channel.queueDeclare(Constants.***TOPIC_QUEUE1** *, true, false, false, null);
channel.queueDeclare(Constants.***TOPIC_QUEUE2** *, true, false, false, null);
*// 4. 声明交换机*
* *channel.exchangeDeclare(Constants.***TOPIC_EXCHANGE** *, BuiltinExchangeType.***TOPIC** *, true, false, false, null);
***// 5. 绑定交换机和队列** *
* // 队列1绑定error信息*
* *channel.queueBind(Constants.*TOPIC_QUEUE1*, Constants.*TOPIC_EXCHANGE*, "***.error** ");
*// 队列2绑定error和info信息*
* *channel.queueBind(Constants.*TOPIC_QUEUE2*, Constants.*TOPIC_EXCHANGE*, "**#.info** ");
channel.queueBind(Constants.*TOPIC_QUEUE2*, Constants.*TOPIC_EXCHANGE*, "***.error** ");
*// 6. 发送消息*
* *String text1 = "hello topic, i am order.error!";
channel.basicPublish(Constants.*TOPIC_EXCHANGE*, "**order.error** ", null, text1.getBytes(StandardCharsets.*UTF_8*));
String text2 = "hello routing, i am order.pay.info!";
channel.basicPublish(Constants.*TOPIC_EXCHANGE*, "**order.pay.info** ", null, text2.getBytes(StandardCharsets.*UTF_8*));
String text3 = "hello routing, i am pay.error!";
channel.basicPublish(Constants.*TOPIC_EXCHANGE*, "**pay.error** ", null, text3.getBytes(StandardCharsets.*UTF_8*));
*// 7. 释放资源*
* *connection.close();
}
}
消费者代码:(有两份,除了绑定队列不同外,基本都是一样的,这里只放一份)
public class consumer1 {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂*
* *ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
*// 2. 创建连接Connection和通道*
* *Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
*// 3. 声明队列(这是安全性措施,因为如果生产者还没创建队列的话,消费者这边直接读取会报错)*
* *channel.queueDeclare(Constants.*TOPIC_QUEUE1*, true, false, false, null);
*// 4. 接收消息,进行消费💥*
* *DefaultConsumer defaultConsumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body) throws IOException
{
System.*out*.println("接收到消息: " + new String(body));
}
};
channel.basicConsume(Constants.*TOPIC_QUEUE1*, true, defaultConsumer);
}
}
五、RPC(RPC通信模式)

RPC(Remote Procedure Call,远程过程调用)是一种通过网络从远程计算机上请求服务,而不需要了解底层网络的技术,类似于Http远程调用。
RabbitMQ 实现 RPC 通信的过程,大概是通过两个队列实现一个可回调的过程。

编写客户端代码
-
发送请求到请求队列中(需要设置
replyTo以及correlationId) -
接收响应队列中的消息,判断
correlationId是否一致,将一致的消息放到阻塞队列中,以便同步获取 下面是需要用到的常量:
*// rpc模式*
public static String *RPC_REQUEST_QUEUE *= "rpc.request.queue";
public static String *RPC_RESPONSE_QUEUE *= "rpc.response.queue";客户端代码如下所示:
public class client {
private final static BlockingQueue<String> *bq *= new ArrayBlockingQueue<>(1);
public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
*// 1. 创建连接工厂、连接Connection和通道*
* *ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
*// 2. 声明队列*
* *channel.queueDeclare(Constants.***RPC_REQUEST_QUEUE** *, true, false, false, null);
channel.queueDeclare(Constants.***RPC_RESPONSE_QUEUE** *, true, false, false, null);
** ** ***// 3. 发送请求到请求队列中,需要设置属性** **(使用内置交换机时, routingKey要和队列名称一样, 才可以路由到对应的队列上去)*
* *String id = UUID.*randomUUID*().toString();
String text = "hello rpc!";
AMQP.BasicProperties props = new AMQP.BasicProperties
.**Builder** ()
.**replyTo** (Constants.*RPC_RESPONSE_QUEUE*) *// 设置回调队列*
* *.**correlationId** (id) *// 唯一标志本次请求 *
* *.**build** ();
channel.basicPublish("", Constants.*RPC_REQUEST_QUEUE*, props, text.getBytes(StandardCharsets.*UTF_8*));
*// 4. 接收响应队列中的消息(需要放到阻塞队列中,保持接收时候的同步)*
* *DefaultConsumer consumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String resp = new String(body);
System.*out*.println("接收到响应队列中消息:" + resp);
***// 校验CorrelationId是否一致** *
* *if(id.equals(properties.getCorrelationId())) {
*bq*.offer(resp);
}
}
};
channel.basicConsume(Constants.*RPC_RESPONSE_QUEUE*, true, consumer);
***// 获取回调的结果** **(因为是从阻塞队列中拿,所以这里会阻塞)*
* *String result = *bq*.take();
System.*out*.println(" [RPCClient] Result:" + result);
*// 释放资源*
* *connection.close();
}
}编写服务端代码
-
接收请求队列中的消息
-
根据消息内容进行响应处理,把应答结果返回到回调队列中
注意事项💥
-
设置同时最多只能获取一个消息
-
如果不设置
basicQos,RabbitMQ 会使用默认的QoS设置,其prefetchCount默认值为0。当prefetchCount为0时,RabbitMQ 会根据内部实现和当前的网络状况等因素,可能会同时发送多条消息给消费者。这意味着在默认情况下,消费者可能会同时接收到多条消息,但具体数量不是严格保证的,可能会有所波动。 -
在 RPC 模式下,通常期望的是一对一的消息处理,即一个请求对应一个响应。消费者在处理完一个消息并确认之后,才会接收到下一条消息。
// 设置同时最多只能获取一个消息 channel.basicQos(1);
-
-
RabbitMQ消息确定机制
- 在 RabbitMQ 中,
basicConsume方法的autoAck参数用于指定消费者是否应该自动向消息队列确认消息:-
自动确认 (autoAck=true):消息队列在将消息发送给消费者后,会立即从内存中删除该消息。这意味着,如果消费者处理消息失败,消息将丢失,因为消息队列认为消息已经被成功消费。
-
手动确认 (autoAck=false):消息队列在将消息发送给消费者后,需要消费者显式地调用
basicAck方法来确认消息。手动确认提供了更高的可靠性,确保消息不会被意外丢失,适用于消息处理重要且需要确保每个消息都被正确处理的场景。
-
- 在 RabbitMQ 中,
public class server {
public static void main(String[] args) throws IOException, TimeoutException {
*// 1. 创建连接工厂、连接Connection和通道*
* *ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
*// 2. 接收请求*
* // 设置同时最多只能获取一个消息💥*
* ***channel.basicQos(1);**
DefaultConsumer consumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String text = new String(body);
System.*out*.println("接收到请求:" + text);
String resp = "[response] " + text; // 处理业务。这里简单拼接字符串即可
*// 3. 发送响应*
* // 设置属性,将消息放到响应队列中*
* *AMQP.BasicProperties props = new AMQP.BasicProperties
.Builder()
.**correlationId** (properties.getCorrelationId())
.build();
channel.basicPublish("", **properties.getReplyTo()** , props, resp.getBytes(StandardCharsets.*UTF_8*));
*// 因为设置了autoAck=false,所以需要手动ack确定一下💥*
* ***channel.basicAck(envelope.getDeliveryTag(), false);**
}
};
channel.basicConsume(Constants.*RPC_REQUEST_QUEUE*, false, consumer);
}
}
六、Publisher Confirms(发布确认机制)
作为消息中间件,都会面临消息丢失的问题。消息丢失大概分为三种情况:
-
生产者问题 。因为应用程序故障,网络抖动等各种原因,生产者没有成功向 Broker 发送消息。
-
消息中间件自身问题 。生产者成功发送给了 Broker,但是 Broker 没有把消息保存好,导致消息丢失。
-
消费者问题 。Broker 发送消息到消费者,消费者在消费消息时,因为没有处理好,导致 Broker 将消费失败的消息从队列中删除了。

RabbitMQ 也对上述问题给出了相应的解决方案。
-
针对问题1,采用 发布确认机制 解决
-
针对问题2,采用 持久化机制 解决
-
针对问题3,采用 消息应答机制 解决
前面一直使用的 basicPublish() 只是把消息写入到 TCP 缓冲区 ,并不代表消息真的到达了 RabbitMQ 服务器或被持久化。
在 Publisher Confirms 模式 下,只有 waitForConfirms() 或 waitForConfirmsOrDie() 收到确认后,消息才算真正安全投递成功。
| 方法 | 是否阻塞 | 粒度 | 性能 | 异常处理 | 场景 |
|---|---|---|---|---|---|
| basicPublish() | 否 | 不确认 | 高 | 无法检测失败 | 不关心可靠性时 |
| waitForConfirms() | 是 | 单条 | 低 | 返回 false 或超时 | 高安全但低吞吐 |
| waitForConfirmsOrDie() | 是 | 批量 | 中 | 抛异常 | 批量发送 |
| addConfirmListener() | 否 | 异步 | 高 | 回调处理 | 高吞吐系统 |
发布确认是 AMQP 0.9.1 协议的扩展,默认情况下它不会被启用 。生产者通过 channel.confirmSelect() 将信道设置为 confirm 模式。
整体代码框架:
*// 常量类:发布确认机制*
public static String *PUBLISH_CONFIRM_QUEUE1 *= "publish.confirm.queue1";
public static String *PUBLISH_CONFIRM_QUEUE2 *= "publish.confirm.queue2";
public static String *PUBLISH_CONFIRM_QUEUE3 *= "publish.confirm.queue3";
public class PublisherConfirms {
public static final Integer *MESSAGE_SIZE *= 10000; // 发送消息的数量
public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
*// Strategy 1: Publishing Messages Individually*
* PublishingMessagesIndividually*();
*// Strategy 2: Publishing Messages in Batches*
* PublishingMessagesInBatchesy*();
*// Strategy 3: Handling Publisher Confirms Asynchronously*
* HandlingPublisherConfirmsAsynchronously*();
}
// 获取rabbitmq连接
public static Connection getConnection() throws IOException, TimeoutException {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost(Constants.*IP*);
factory.setPort(Constants.*PORT*);
factory.setVirtualHost(Constants.*VIRTUALHOST*);
factory.setUsername(Constants.*USERNAME*);
factory.setPassword(Constants.*PASSWORD*);
Connection connection = factory.newConnection();
return connection;
}
*// 单独确认*
* *public static void PublishingMessagesIndividually() {}
*// 批量确认*
* *public static void PublishingMessagesInBatchesy() {}
*// 异步确认*
* *public static void HandlingPublisherConfirmsAsynchronously() {}
}① Publishing Messages Individually(单条确认)
这种策略是每发送一条消息后就调用 channel.waitForConfirmsOrDie()方法 ,之后等待服务端的确认,这实际上是一种串行同步等待 的方式。尤其对于持久化的消息来说,需要等待消息确认存储在磁盘之后才会返回(调用Linux内核的fsync方法)
*// 单独确认
public static void PublishingMessagesIndividually() throws IOException, TimeoutException, InterruptedException {
try (Connection connection = getConnection()) {
// 1. 获取通道
Channel channel = connection.createChannel();
// 2. 开启confirm模式
** ****channel.confirmSelect();** **
// 3. 声明队列
channel.queueDeclare(Constants.PUBLISH_CONFIRM_QUEUE1, true, false, true, null);
// 4. 发送消息
long start = System.currentTimeMillis();
for(int i = 0; i < MESSAGE_SIZE; ++i) {
String text = "hello publisher_confirms";
channel.basicPublish("", Constants.PUBLISH_CONFIRM_QUEUE1, null, text.getBytes(StandardCharsets.UTF_8));
// waitForConfirmsOrDie() 阻塞当前线程,直到 RabbitMQ 返回所有已发送消息的确认。
// 如果超时过期, 则抛出TimeoutException。如果任何消息被nack(丢失), 则抛出IOException。
** ****channel.waitForConfirmsOrDie(5000);** **
}
long end = System.currentTimeMillis();
System.out.println("PublishingMessagesIndividually花费了:" + (end - start) + "ms");
}
**}*② Publishing Messages in Batches(批量确认)
相比于单独确认策略,批量确认可以一次性发送多条消息,再批量进行消息确认,极大地提升效率!
缺点是出现 Basic.Nack 或者超时的情况,我们不清楚具体哪条消息出了问题,客户端需要将这一批次的消息全部重发,这会带来明显的重复消息数量。当消息经常丢失时,批量确认的性能应该是不升反降的 。
*// 批量确认
public static void PublishingMessagesInBatchesy() throws IOException, TimeoutException, InterruptedException {
try (Connection connection = getConnection()) {
// 1. 获取通道
Channel channel = connection.createChannel();
// 2. 开启confirm模式
** ****channel.confirmSelect();** **
// 3. 声明队列
channel.queueDeclare(Constants.PUBLISH_CONFIRM_QUEUE2, true, false, true, null);
// 4. 发送消息
** int ****batchSize** ** = 200; // 一次确认的消息数量
** int ****outstandingMessageCount** ** = 0; // 记录当前已经
long start = System.currentTimeMillis();
for (int i = 0; i < MESSAGE_SIZE; ++i) {
// basicPublish() 只是把消息写入到 TCP 缓冲区,并不代表消息真的到达了 RabbitMQ 服务器或被持久化。
// 💥在 Publisher Confirms 模式下,只有 waitForConfirms() 或 waitForConfirmsOrDie() 收到确认后,消息才算真正安全投递成功。
String text = "hello publisher_confirms";
channel.basicPublish("", Constants.PUBLISH_CONFIRM_QUEUE2, null, text.getBytes(StandardCharsets.UTF_8));
** ****outstandingMessageCount++;** **
** ****// 批量确认消息** **
if (outstandingMessageCount == batchSize) {
channel.waitForConfirmsOrDie(5000);
outstandingMessageCount = 0;
}
}
** ****// 消息发送完, 还有未确认的消息, 则进行确认** **
if (outstandingMessageCount > 0) {
channel.waitForConfirmsOrDie(5000);
}
long end = System.currentTimeMillis();
System.out.println("PublishingMessagesInBatchesy花费了:" + (end - start) + "ms");
** }*③ Handling Publisher Confirms Asynchronously(异步确认)
生产者将信道设置成 confirm模式 ,一旦信道进入 confirm 模式,所有在该信道上面发布的消息都会被指派一个唯一的ID(从1开始),一旦消息被投递到所有匹配的队列之后,RabbitMQ 就会发送一个确认给生产者(包含消息的唯一ID),这就使得生产者知道消息已经正确到达目的队列了,如果消息和队列是可持久化的,那么确认消息 ack 会在将消息写入磁盘之后发出。
Broker 回传给生产者的确认消息中 deliveryTag包含了确认消息的序号 ,此外 Broker 也可以设置 channel.basicAck 方法中的 multiple 参数,表示到这个序号之前的所有消息都已经得到了处理。

Channel 接口提供了一个方法 addConfirmListener,这个方法可以添加 ConfirmListener 接口,这个接口中包含两个方法,分别对应处理 RabbitMQ 发送给生产者的 ack 和 nack:
-
handleAck``(long deliveryTag, boolean multiple) -
handleNack``(long deliveryTag, boolean multiple)-
deliveryTag:表示发送消息的序号 -
multiple:表示是否为批量确认
-
此外,在编写代码的时候,需要为每一个 Channel 维护一个已发送消息的序号集合,当收到 RabbitMQ 的 confirm 回调时,从集合中删除对应确认的消息。当 Channel 开启 confirm 模式后,Channel 上发送消息都会附带一个从 1 开始递增的 deliveryTag 序号,我们可以使用 SortedSet的有序性来维护这个已发消息的集合 。
-
当收到 ack 时,从序列中删除该消息的序号。如果为批量确认消息,表示小于等于当前序号 deliveryTag 的消息都收到了,则清除对应集合。
-
当收到 nack 时,处理逻辑类似,不过需要结合具体的业务情况,进行消息重发等操作。
*// 异步确认*
public static void HandlingPublisherConfirmsAsynchronously() throws IOException, TimeoutException, InterruptedException {
try (Connection connection = *getConnection*()) {
*// 1. 获取通道*
* *Channel channel = connection.createChannel();
*// 2. 开启confirm模式*
* ***channel.confirmSelect();**
*// 3. 声明队列*
* *channel.queueDeclare(Constants.*PUBLISH_CONFIRM_QUEUE3*, true, false, true, null);
*// 创建一个有序集合SortedSet,存放delivery序号*
* *SortedSet<Long> set = Collections.***synchronizedSortedSet** *(new TreeSet<>());
*// 4. 添加回调接口*
* ***channel.addConfirmListener** (
(deliveryTag, multiple) -> {
if (multiple) {
*// 批量确认:获取小于等于deliveryTag的序号集合,进行删除,表示这批序号的消息都已经被ack了*
* ***set.headSet(deliveryTag + 1).clear()** ;
} else {
*// 单条确认:将当前的deliveryTag从集合中移除*
* ***set.remove(deliveryTag)** ;
}
},
(deliveryTag, multiple) -> {
if (multiple) {
*// 批量确认:获取小于等于deliveryTag的序号集合,进行删除,表示这批序号的消息都已经被ack了*
* *set.headSet(deliveryTag + 1).clear();
} else {
*// 单条确认:将当前的deliveryTag从集合中移除*
* *set.remove(deliveryTag);
}
*// 如果处理失败, 这里需要添加处理消息重发的场景,此处代码省略*
* *}
);
*// 5. 发送消息*
* *long start = System.*currentTimeMillis*();
for(int i = 0; i < *MESSAGE_SIZE*; ++i) {
String text = "hello publisher_confirms";
*// 获取下一次发送的序号,必须在basicPublish之前调用,否则会出现错位!💥*
* *long nextPublishSeqNo = **channel.getNextPublishSeqNo()** ;
channel.basicPublish("", Constants.*PUBLISH_CONFIRM_QUEUE3*, null, text.getBytes(StandardCharsets.*UTF_8*));
*// 将序号存放到有序集合中*
* ***set.add(nextPublishSeqNo)** ;
}
*// 确认消息都确认完毕*
* *while(!set.isEmpty()) {
Thread.*sleep*(10);
}
long end = System.*currentTimeMillis*();
System.*out*.println("PublishingMessagesInBatchesy花费了:" + (end - start) + "ms");
}
}运行结果如下所示:
PublishingMessagesIndividually花费了:2738ms
PublishingMessagesInBatchesy花费了:352ms
HandlingPublisherConfirmsAsynchronously花费了:192msⅢ. SpringBoot整合RabbitMQ
Spring 官方:Spring AMQP
RabbitMQ 官方:RabbitMQ tutorial - "Hello World!" | RabbitMQ
一、Work Queue(工作队列模式)

步骤: (后面其它模式也是如此)
-
引入依赖
-
编写 yml 配置文件,基本信息配置
-
编写生产者代码
-
编写消费者代码
- 定义监听类,使用
@RabbitListener注解完成队列监听
- 定义监听类,使用
-
运行观察结果
引入依赖
<!--Spring MVC相关依赖-->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<!--RabbitMQ相关依赖-->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>添加配置
# 配置RabbitMQ的基本信息
spring:
rabbitmq:
host: 127.0.0.1
port: 5672
username: liren
password: 123123
virtual-host: lirendada编写生产者代码
常量类:
public class Constants {
public static final String *WORK_QUEUE *= "work_queue";
}然后在 config 包中声明队列:(注意包要导对~)
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.QueueBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMQConfig {
@Bean("workQueue")
public Queue WorkQueue() {
return **QueueBuilder** .*durable*(Constants.*WORK_QUEUE*).**build** ();
}
}最后在需要发送消息的地方调用 RabbitTemplate 发送消息:
@RequestMapping("/producer")
@RestController
public class ProducerController {
@Autowired
**private RabbitTemplate rabbitTemplate** ;
@RequestMapping("/work")
public String work(){
for (int i = 0; i < 10; i++) {
// 使用内置交换机发送消息, routingKey和队列名称保持一致
rabbitTemplate.**convertAndSend** ("", Constants.WORK_QUEUE, "hello spring amqp: work...");
}
return "发送成功";
}
}编写消费者代码
定义监听类,用于消费队列中的消息:(注意包要导对~)
import com.liren.springbootrabbitmq.constant.Constants;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
public class WorkListener {
**@RabbitListener(queues = Constants.** ***WORK_QUEUE** ***)**
public void workqueue1(Message message) {
System.*out*.println("workqueue1 [" + Constants.*WORK_QUEUE *+ "]收到消息:" + message);
}
**@RabbitListener(queues = Constants.** ***WORK_QUEUE** ***)**
public void workqueue2(String message) {
System.*out*.println("workqueue2 [" + Constants.*WORK_QUEUE *+ "]收到消息:" + message);
}
}@RabbitListener 是 Spring 框架中用于监听 RabbitMQ 队列的注解,通过使用这个注解,可以定义一个方法,以便从 RabbitMQ 队列中接收消息。该注解支持多种参数类型,这些参数类型代表了从 RabbitMQ 接收到的消息和相关信息。
以下是一些常用的参数类型:
-
String:返回消息的内容 -
Message(org.springframework.amqp.core.Message):Spring AMQP 的Message类,返回原始的消息体以及消息的属性,如消息ID、内容、队列信息等。 -
Channel(com.rabbitmq.client.Channel):RabbitMQ 的通道对象,可以用于进行更高级的操作,如手动确认消息。
运行结果
运行程序,然后发起请求,会有三个队列接收消息,如下所示:

管理页面中可以看到三个消费者以及一个生产者通道:

二、Publish/Subscribe(发布/订阅模式)

RabbitMQ 交换机常见三种类型:fanout、direct、topic,不同类型有着不同的路由策略。
-
Fanout :广播策略 ,将消息交给所有绑定到该交换机的队列(Publish/Subscribe 模式 )
-
Direct :定向策略 ,把消息交给符合指定
routing key的队列(Routing 模式 ) -
Topic :通配符策略 ,把消息交给符合
routing pattern的队列(Topics 模式 )
编写生产者代码
常量类:
*// 发布订阅模式*
public static final String *FANOUT_QUEUE1 *= "fanout.queue1";
public static final String *FANOUT_QUEUE2 *= "fanout.queue2";
public static final String *FANOUT_EXCHANGE *= "fanout.exchange";然后在 config 包中声明队列:(注意包要导对~)
*// 发布订阅模式*
@Bean("publishConfirmQueue1")
public Queue publishConfirmQueue1() {
return QueueBuilder.*durable*(Constants.*FANOUT_QUEUE1*).build(); *// 声明队列*
}
@Bean("publishConfirmQueue2")
public Queue publishConfirmQueue2() {
return QueueBuilder.*durable*(Constants.*FANOUT_QUEUE2*).build(); *// 声明队列*
}
@Bean("fanoutExchange")
public FanoutExchange fanoutExchange() {
return **ExchangeBuilder** .***fanoutExchange** *(Constants.*FANOUT_EXCHANGE*).build(); *// 声明交换机*
}
@Bean("fanoutBinding1")
public Binding fanoutBinding1(@Qualifier("publishConfirmQueue1") Queue queue,
@Qualifier("fanoutExchange") FanoutExchange exchange) {
return **BindingBuilder** .***bind** *(queue).**to** (exchange); *// 绑定交换机和队列*
}
@Bean("fanoutBinding2")
public Binding fanoutBinding2(@Qualifier("publishConfirmQueue2") Queue queue,
@Qualifier("fanoutExchange") FanoutExchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange); *// 绑定交换机和队列*
}使用接口发送消息
@RequestMapping("/fanout")
public String fanout() {
rabbitTemplate.convertAndSend(**Constants.** ***FANOUT_EXCHANGE** *, **""** , "hello spring amqp: fanout...");
return "发送成功!";
}编写消费者代码
@Component
public class FanoutListener {
@RabbitListener(queues = Constants.*FANOUT_QUEUE1*)
public void fanoutQueue1(String message) {
System.*out*.println("fanoutQueue1 [" + Constants.*FANOUT_QUEUE1 *+ "]收到消息:" + message);
}
@RabbitListener(queues = Constants.*FANOUT_QUEUE2*)
public void fanoutQueue2(String message) {
System.*out*.println("fanoutQueue2 [" + Constants.*FANOUT_QUEUE2 *+ "]收到消息:" + message);
}
}
消费者另一种写法
@RabbitListener 是一个功能强大的注解。这个注解里面可以配置 @QueueBinding、@Queue、@Exchange,直接通过这个组合注解一次性搞定多个交换机、绑定、路由、并且配置监听功能等
@Slf4j
@Component
public class UserRegisterListener {
@RabbitListener(
**bindings = @QueueBinding** (
value = **@Queue** (
value = Constants.USER_QUEUE_NANE, // 队列名
durable = "true" // 是否持久化
),
exchange = **@Exchange** (
value = Constants.USER_EXCHANGE_NAME, // 交换机名
type = ExchangeTypes.FANOUT // fanout 交换机
)
// fanout 不需要 routingKey
)
)
public void MailListenerQueue(Message message, Channel channel) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
*// 处理用户注册消息*
* *String body = new String(message.getBody());
*log*.info("用户注册消息处理成功,deliveryTag={}, message={}", deliveryTag, body);
*// 发送邮件TODO*
* *
* // 确认消息*
* *channel.basicAck(deliveryTag, true);
}catch (Exception e) {
*// 异常拒绝消息,进行重发*
* *channel.basicNack(deliveryTag, true, true);
*log*.error("用户注册消息处理失败,拒绝消息,deliveryTag={}", deliveryTag, e);
}
}
}启动时 Spring AMQP 会做的事情,顺序大致是:
-
QueueDeclare- 声明一个 durable 队列
-
ExchangeDeclare- 声明一个 fanout 交换机
-
QueueBind- 把队列绑定到交换机
三、Routing(路由模式)

RabbitMQ 交换机常见三种类型:fanout、direct、topic,不同类型有着不同的路由策略。
-
Fanout :广播策略 ,将消息交给所有绑定到该交换机的队列(Publish/Subscribe 模式 )
-
Direct :定向策略 ,把消息交给符合指定
routing key的队列(Routing 模式 ) -
Topic :通配符策略 ,把消息交给符合
routing pattern的队列(Topics 模式 )
路由模式采用的是 RabbitMQ 中的 Direct 定向策略,生产者发送消息的时候,交换机需要根据消息中的 Routing Key将消息发送给指定的队列 ,而不是发给每一个队列了!
此时,队列和交换机的绑定,不能是任意的绑定了,而是要指定一个 Binding Key。
只有队列绑定时的 Binding Key和消息中的 Routing Key完全一致,队列才会接收到消息 。
编写生产者代码
常量类:
*// 路由模式
public static final String DIRECT_EXCHANGE = "direct.exchange";
public static final String DIRECT_QUEUE1 = "direct.queue1";
**public static final String DIRECT_QUEUE2 = "direct.queue2";*和发布订阅模式的区别是:交换机类型不同、绑定队列的 Binding Key 不同。
*// 路由模式(direct模式)*
@Bean("directQueue1")
public Queue directQueue1() {
return QueueBuilder.*durable*(Constants.*DIRECT_QUEUE1*).build(); *// 声明队列*
}
@Bean("directQueue2")
public Queue directQueue2() {
return QueueBuilder.*durable*(Constants.*DIRECT_QUEUE2*).build(); *// 声明队列*
}
@Bean("directExchange")
public DirectExchange directExchange() {
return ExchangeBuilder.*directExchange*(Constants.*DIRECT_EXCHANGE*).build(); *// 声明交换机*
}
*// 队列1绑定orange*
@Bean("directBinding1")
public Binding directBinding1(@Qualifier("directQueue1") Queue queue,
@Qualifier("directExchange") DirectExchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).**with("orange")** ; *// 绑定交换机和队列,以及bindingKey*
}
*// 队列2绑定green、black*
@Bean("directBinding2")
public Binding directBinding2(@Qualifier("directQueue2") Queue queue,
@Qualifier("directExchange") DirectExchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).**with("green")** ; *// 绑定交换机和队列,以及bindingKey*
}
@Bean("directBinding3")
public Binding directBinding3(@Qualifier("directQueue2") Queue queue,
@Qualifier("directExchange") DirectExchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).**with("black")** ; *// 绑定交换机和队列,以及bindingKey*
}使用接口发送消息:
@RequestMapping("/direct/{routing_key}")
public String dirct(@PathVariable("routing_key") String routing_key) {
rabbitTemplate.convertAndSend(Constants.*DIRECT_EXCHANGE*, routing_key, "hello spring amqp: direct..." + routing_key);
return "发送成功!";
}编写消费者代码
@Component
public class DirectListener {
@RabbitListener(queues = Constants.*DIRECT_QUEUE1*)
public void directQueue1(String message) {
System.*out*.println("directQueue1 [" + Constants.*DIRECT_QUEUE1 *+ "]收到消息:" + message);
}
@RabbitListener(queues = Constants.*DIRECT_QUEUE2*)
public void directQueue2(String message) {
System.*out*.println("directQueue2 [" + Constants.*DIRECT_QUEUE2 *+ "]收到消息:" + message);
}
}分别请求三个不同的 routingkey,结果如下所示:

四、Topics(通配符模式)

Topics 和 Routing 模式的区别是:
-
交换机类型不同: Topics 模式使用的交换机类型为
topic;Routing 模式使用的交换机类型为direct。 -
匹配规则不同:
topic类型的交换机在匹配规则上进行了扩展,Binding Key支持通配符匹配;direct类型的交换机路由规则是Binding Key和Routing Key完全匹配。
匹配规则有如下要求:
-
Routing Key是一系列由点.分隔的单词 ,比如 "stock.usd.nyse"、"nyse.vmw"、"quick.orange.rabbit" -
Binding Key和Routing Key一样,也是点.分割的字符串 -
Binding Key中可以存在两种特殊字符串,用于模糊匹配-
*:表示一个单词 -
#:表示多个单词(0-N个)
-
编写生产者代码
常量类:
*// 通配符模式*
public static final String *TOPIC_EXCHANGE *= "topic.exchange";
public static final String *TOPIC_QUEUE1 *= "topic.queue1";
public static final String *TOPIC_QUEUE2 *= "topic.queue2";生产者代码如下所示:
*// 通配符模式(topics模式)*
@Bean("topicQueue1")
public Queue topicQueue1() {
return QueueBuilder.*durable*(Constants.*TOPIC_QUEUE1*).build(); *// 声明队列*
}
@Bean("topicQueue2")
public Queue topicQueue2() {
return QueueBuilder.*durable*(Constants.*TOPIC_QUEUE2*).build(); *// 声明队列*
}
@Bean("topicExchange")
public **TopicExchange** topicExchange() {
return ExchangeBuilder.*topicExchange*(Constants.*TOPIC_EXCHANGE*).build(); *// 声明交换机*
}
*// 队列1绑定error, 仅接收error信息*
@Bean("topicBinding1")
public Binding topicBinding1(@Qualifier("topicQueue1") Queue queue,
@Qualifier("topicExchange") TopicExchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).**with("*.error")** ; *// 绑定交换机和队列,以及bindingKey*
}
*// 队列2绑定info, error: error,info信息都接收*
@Bean("topicBinding2")
public Binding topicBinding2(@Qualifier("topicQueue2") Queue queue,
@Qualifier("topicExchange") TopicExchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).**with("*.error")** ; *// 绑定交换机和队列,以及bindingKey*
}
@Bean("topicBinding3")
public Binding topicBinding3(@Qualifier("topicQueue2") Queue queue,
@Qualifier("topicExchange") TopicExchange exchange) {
return BindingBuilder.*bind*(queue).to(exchange).**with("#.info")** ; *// 绑定交换机和队列,以及bindingKey*
}使用接口发送消息:
@RequestMapping("/topics/{routing_key}")
public String topics(@PathVariable("routing_key") String routing_key) {
rabbitTemplate.convertAndSend(Constants.*TOPIC_EXCHANGE*, routing_key, "hello spring amqp: topics..." + routing_key);
return "发送成功!";
}编写消费者代码
@Component
public class TopicListener {
@RabbitListener(queues = Constants.*TOPIC_QUEUE1*)
public void topicQueue1(String message) {
System.*out*.println("topicQueue1 [" + Constants.*TOPIC_QUEUE1 *+ "]收到消息:" + message);
}
@RabbitListener(queues = Constants.*TOPIC_QUEUE2*)
public void topicQueue2(String message) {
System.*out*.println("topicQueue2 [" + Constants.*TOPIC_QUEUE2 *+ "]收到消息:" + message);
}
}分别请求两个不同的请求以及参数之后,运行结果如下:


