赞
踩
目录
在介绍Rabbitmq的五种常用的工作模式之前先导入需要的依赖
<dependencies>
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.14.2</version>
</dependency>
</dependencies>
然后创建两个模块,一个是生产者模块,一个是消费者模块,方便我们后面演示效果
首先介绍一下,简单模式,在上图的模型中,有以下概念:
P:生产者,也就是要发送消息的程序
C:消费者,消息的接收者,会一直等待消息到来
queue:消息队列,图中红色部分。类似一个邮箱,可以缓存消息;生产者向其中投递消息,消费者从其中取出消息
这里需要注意的是,生产者和消费者不在一个模块中
- public class HelloProducer {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
- //创建队列
- /**
- * String queue, 队列的名称. 如果该名称不存在 则创建如果存在则不创建
- * boolean durable, 该对象是否持久化 当rabbitmq重启后 队列就会消失
- * boolean exclusive, 该队列是否被一个消费者独占
- * boolean autoDelete,当没有消费者时,该队列是否被自动删除
- * Map<String, Object> arguments: 额外参数的设置
- */
- channel.queueDeclare("hello_queue", true, false, false, null);
- //发送消息
- /**
- * String exchange, 交换机的名称 简单模式没有交换机使用""表示采用默认交换机
- * String routingKey, 路由标识 如果是简单模式起名为队列的名称
- * BasicProperties props, 消息的属性设置。设置为null
- * byte[] body: 消息的内容
- */
- //msg是测试发送的信息
- String msg="hello rabbitmq";
- channel.basicPublish("","hello_queue",null,msg.getBytes());
- //关闭资源
- channel.close();
- connection. Close();
- }
- }
- public class HelloConsumer {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- System.out.println("消费者的标志:"+consumerTag);
- System.out.println("交换机名称:"+envelope.getExchange());
- System.out.println("路由key标志:"+envelope.getRoutingKey());
- System.out.println("消息属性:"+properties);
- super.handleDelivery(consumerTag, envelope, properties, body);
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("hello_queue",true,consumer);
-
- //这里需要注意一点,不能关闭connection和channel
- }
- }
首先运行HelloProducer项目,然后在RabbitMQ前端控制台可以看到队列中此时存在一个Messages,因为此时没有消费者(消费者没有启动),所以此时消息存在于队列中
启动HelloConsumer ,控制台打印结果
此时观察RabbitMQ前端控制器队列中的情况
Work Queues:与入门程序的简单模式相比,多了一个或一些消费端,多个消费端共同消费同一个队列中的消息。
应用场景:对于任务过重或任务较多情况使用工作队列可以提高任务处理的速度。
- public class WorkProducer {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
- //创建队列
- /**
- * String queue, 队列的名称. 如果该名称不存在 则创建如果存在则不创建
- * boolean durable, 该对象是否持久化 当rabbitmq重启后 队列就会消失
- * boolean exclusive, 该队列是否被一个消费者独占
- * boolean autoDelete,当没有消费者时,该队列是否被自动删除
- * Map<String, Object> arguments: 额外参数的设置
- */
- channel.queueDeclare("work_queue", true, false, false, null);
- //发送消息
- /**
- * String exchange, 交换机的名称 简单模式没有交换机使用""表示采用默认交换机
- * String routingKey, 路由标识 如果是简单模式起名为队列的名称
- * BasicProperties props, 消息的属性设置。设置为null
- * byte[] body: 消息的内容
- */
- //msg是测试发送的信息(这里和simple模式不一样,我们一次性生产十条消息)
- for (int i = 0; i < 10; i++) {
- String msg="work rabbitmq-----------------"+i;
- channel.basicPublish("","work_queue",null,msg.getBytes());
- }
- //关闭资源
- channel.close();
- connection. Close();
- }
- }
- //消费者1
- public class WorkConsumer01 {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("hello_queue",true,consumer);
-
- //这里需要注意一点,不能关闭connection和channel
- }
- }
-
-
- //消费者2
- public class WorkConsumer02 {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("hello_queue",true,consumer);
-
- //这里需要注意一点,不能关闭connection和channel
- }
- }
首先运行生产者,接着运行两个消费者,然后观察idea控制台打印结果
消费者1
消费者2
在一个队列中如果有多个消费者,那么消费者之间对于同一个消息的关系是竞争的关系。Work Queues 对于任务过重或任务较多情况使用工作队列可以提高任务处理的速度。例如:短信服务部署多个,只需要有一个节点成功发送即可。
在订阅模型中,多了一个 Exchange 角色,而且过程略有变化:
P:生产者,也就是要发送消息的程序,但是不再发送到队列中,而是发给X(交换机)
C:消费者,消息的接收者,会一直等待消息到来
Queue:消息队列,接收消息、缓存消息
Exchange:交换机(X)。一方面,接收生产者发送的消息。另一方面,知道如何处理消息,例如递交给某个特别队列、递交给所有队列、或是将消息丢弃。到底如何操作,取决于Exchange的类型。Exchange有常见以下3种类型:
1. Fanout:广播,将消息交给所有绑定到交换机的队列
2. Direct:定向,把消息交给符合指定routing key 的队列
3. Topic:通配符,把消息交给符合routing pattern(路由模式) 的队列
Exchange(交换机)只负责转发消息,不具备存储消息的能力,因此如果没有任何队列与Exchange 绑定,或者没有符合路由规则的队列,那么消息会丢失!
- public class PublishProducer {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
- //创建交换机
- /**
- * String exchange, 交换机的名称 如果不存在则创建,存在则不创建
- * BuiltinExchangeType type, 交换机的类型
- * boolean durable: 是否持久化。
- */
- channel.exchangeDeclare("publish_exchange", BuiltinExchangeType.FANOUT,true);
- //创建队列
- /**
- * String queue, 队列的名称. 如果该名称不存在 则创建如果存在则不创建
- * boolean durable, 该对象是否持久化 当rabbitmq重启后 队列就会消失
- * boolean exclusive, 该队列是否被一个消费者独占
- * boolean autoDelete,当没有消费者时,该队列是否被自动删除
- * Map<String, Object> arguments: 额外参数的设置
- */
- //队列1
- channel.queueDeclare("publish_queue01", true, false, false, null);
- //队列2
- channel.queueDeclare("publish_queue02", true, false, false, null);
- //将队列和交换机绑定
- /**
- * String queue
- * String exchange
- * String routingKey:发布订阅模式没有routingKey则写为""
- */
- channel.queueBind("publish_queue01","publish_exchange","");
- channel.queueBind("publish_queue02","publish_exchange","");
- //发送消息
- /**
- * String exchange, 交换机的名称 简单模式没有交换机使用""表示采用默认交换机
- * String routingKey, 路由标识 如果是简单模式起名为队列的名称
- * BasicProperties props, 消息的属性设置。设置为null
- * byte[] body: 消息的内容
- */
- String msg="publish_exchange rabbitmq";
- channel.basicPublish("publish_exchange","",null,msg.getBytes());
- //关闭资源
- channel.close();
- connection. Close();
- }
- }
- //消费者1
- public class PublishConsumer01 {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
-
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("publish_queue01",true,consumer);
-
- }
- }
-
- //消费者2
- public class PublishConsumer02 {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
-
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("publish_queue02",true,consumer);
-
- }
- }
首先运行生产者,接着运行两个消费者,然后观察idea控制台打印结果
消费者1
消费者2
交换机需要与队列进行绑定,绑定之后;一个消息可以被多个消费者都收到。
2. 发布订阅模式与工作队列模式的区别:
工作队列模式不用定义交换机,而发布/订阅模式需要定义交换机发布/订阅模式的生产方是面向交换机发送消息,工作队列模式的生产方是面向队列发送消息(底层使用默认交换机)
发布/订阅模式需要设置队列和交换机的绑定,工作队列模式不需要设置,实际上工作队列模式会将队列绑 定到默认的交换机
队列与交换机的绑定,不能是任意绑定了,而是要指定一个RoutingKey(路由key),消息的发送方在向 Exchange 发送消息时,也必须指定消息的RoutingKey
Exchange 不再把消息交给每一个绑定的队列,而是根据消息的Routing Key 进行判断,只有队列的Routingkey 与消息的 Routingkey 完全一致,才会接收到消息
P:生产者,向 Exchange 发送消息,发送消息时,会指定一个routing key
X:Exchange(交换机),接收生产者的消息,然后把消息递交给与 routing key 完全匹配的队列
C1:消费者,其所在队列指定了需要 routing key 为 error 的消息
C2:消费者,其所在队列指定了需要 routing key 为 info、error、warning 的消息
- public class RoutingProducer {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
- //创建交换机
- /**
- * String exchange, 交换机的名称 如果不存在则创建,存在则不创建
- * BuiltinExchangeType type, 交换机的类型
- * boolean durable: 是否持久化。
- */
- channel.exchangeDeclare("router_exchange", BuiltinExchangeType.DIRECT,true);
- //创建队列
- /**
- * String queue, 队列的名称. 如果该名称不存在 则创建如果存在则不创建
- * boolean durable, 该对象是否持久化 当rabbitmq重启后 队列就会消失
- * boolean exclusive, 该队列是否被一个消费者独占
- * boolean autoDelete,当没有消费者时,该队列是否被自动删除
- * Map<String, Object> arguments: 额外参数的设置
- */
- //队列1
- channel.queueDeclare("router_queue01", true, false, false, null);
- //队列2
- channel.queueDeclare("router_queue02", true, false, false, null);
- //将队列和交换机绑定
- /**
- * String queue
- * String exchange
- * String routingKey:发布订阅模式没有routingKey则写为""
- */
- channel.queueBind("router_queue01","router_exchange","error");
- channel.queueBind("router_queue02","router_exchange","error");
- channel.queueBind("router_queue02","router_exchange","info");
- channel.queueBind("router_queue02","router_exchange","warning");
- //发送消息
- /**
- * String exchange, 交换机的名称 简单模式没有交换机使用""表示采用默认交换机
- * String routingKey, 路由标识 如果是简单模式起名为队列的名称
- * BasicProperties props, 消息的属性设置。设置为null
- * byte[] body: 消息的内容
- */
- String msg="router_exchange rabbitmq";
- channel.basicPublish("router_exchange","info",null,msg.getBytes());
- //关闭资源
- channel.close();
- connection.close();
- }
- }
- //消费者1
- public class RoutingConsumer01 {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
-
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("router_queue01",true,consumer);
-
- }
- }
- //消费者2
- public class RoutingConsumer02 {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
-
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("router_queue02",true,consumer);
-
- }
- }
首先运行生产者,接着运行两个消费者,然后观察idea控制台打印结果
消费者1
消费者2
Routing 模式要求队列在绑定交换机时要指定 routing key,消息会转发到符合 routing key 的队列。
Topic 类型与 Direct 相比,都是可以根据 RoutingKey 把消息路由到不同的队列。只不过 Topic 类型Exchange 可以让队列在绑定Routing key 的时候使用通配符!
Routingkey 一般都是有一个或多个单词组成,多个单词之间以”.”分割,例如: item.insert
通配符规则:# 匹配一个或多个词,* 匹配不多不少恰好1个词,例如:item.# 能够匹配 item.insert.abc 或者 item.insert,item.*只能匹配 item.insert
- public class TopicsProducer {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
- //创建交换机
- /**
- * String exchange, 交换机的名称 如果不存在则创建,存在则不创建
- * BuiltinExchangeType type, 交换机的类型
- * boolean durable: 是否持久化。
- */
- channel.exchangeDeclare("topic_exchange", BuiltinExchangeType.TOPIC,true);
- //创建队列
- /**
- * String queue, 队列的名称. 如果该名称不存在 则创建如果存在则不创建
- * boolean durable, 该对象是否持久化 当rabbitmq重启后 队列就会消失
- * boolean exclusive, 该队列是否被一个消费者独占
- * boolean autoDelete,当没有消费者时,该队列是否被自动删除
- * Map<String, Object> arguments: 额外参数的设置
- */
- //队列1
- channel.queueDeclare("topic_queue01", true, false, false, null);
- //队列2
- channel.queueDeclare("topic_queue02", true, false, false, null);
- //将队列和交换机绑定
- /**
- * String queue
- * String exchange
- * String routingKey:发布订阅模式没有routingKey则写为""
- */
- channel.queueBind("topic_queue01","topic_exchange","*.orange.*");
- channel.queueBind("topic_queue02","topic_exchange","*.*.rabbit");
- channel.queueBind("topic_queue02","topic_exchange","lazy.#");
- //发送消息
- /**
- * String exchange, 交换机的名称 简单模式没有交换机使用""表示采用默认交换机
- * String routingKey, 路由标识 如果是简单模式起名为队列的名称
- * BasicProperties props, 消息的属性设置。设置为null
- * byte[] body: 消息的内容
- */
- String msg="topic_exchange rabbitmq";
- channel.basicPublish("topic_exchange","lazy.rabbit.orange",null,msg.getBytes());
- //关闭资源
- channel.close();
- connection. Close();
- }
- }
- //消费者1
- public class TopicsConsumer01 {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
-
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("topic_queue01",true,consumer);
-
- }
- }
- //消费者2
- public class TopicsConsumer02 {
- public static void main(String[] args) throws IOException, TimeoutException {
- //创建连接,并设置连接信息
- ConnectionFactory factory = new ConnectionFactory();
- //设置rabbitmq服务的地址
- factory.setHost("192.168.47.135");
- //设置rabbitmq服务的端口号
- factory.setPort(5672);
- //设置连接的账号
- factory.setUsername("ykk");
- //设置连接的密码
- factory.setPassword("ykk");
- //设置虚拟主机
- factory.setVirtualHost("/ykk");
- //获取连接对象
- Connection connection = factory.newConnection();
- //获取channel对象
- Channel channel = connection.createChannel();
-
- //接受队列中的消息
- Consumer consumer=new DefaultConsumer(channel){
- /**
- * @param consumerTag: 消费者的标签
- * @param envelope : 设置 拿到你的交换机 路由key等信息
- * @param properties: 消息的属性对象
- * @param body: 消息的内容
- * @throws IOException
- */
- @Override
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
- //这里我们在控制台打印一下,看下消费者是否接收到了生产者传过来的信息
- System.out.println("接受的内容:"+new String(body));
- }
- };
- /**
- * String queue, 队列名
- * boolean autoAck,是否自动确认。 当rabbitmq把消息发送给消费后,消费端自动确认消息。
- * Consumer callback:回调。 当rabbitmq队列中存在消息则触发该回调
- */
- channel.basicConsume("topic_queue02",true,consumer);
-
- }
- }
首先运行生产者,接着运行两个消费者,然后观察idea控制台打印结果
消费者1有结果,消费者2无结果
Topic 主题模式可以实现 Pub/Sub 发布与订阅模式和 Routing 路由模式的功能,只是 Topic 在配置routing key 的时候可以使用通配符,显得更加灵活。
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。