赞
踩
MQ全称 Message Queue(消息队列),是在消息的传输过程中保存消息的容器。多用于分布式系统之间进行通信。
应用之间的远程调用
加入MQ后应用之间的调用
MQ相当于一个中介,生产方通过MQ与消费方交互,它将应用程序进行解耦合。
系统的耦合性越高,容错性就越低,可维护性就越低。
使用 MQ 使得应用间解耦,提升容错性和可维护性。
将不需要同步处理的并且耗时长的操作由消息队列通知消息接收方进行异步处理。提高了应用程序的响应时间。
一个下单操作耗时:20 + 300 + 300 + 300 = 920ms
用户点击完下单按钮后,需要等待920ms才能得到下单响应,太慢!
用户点击完下单按钮后,只需等待25ms就能得到下单响应 (20 + 5 = 25ms)。
提升用户体验和系统吞吐量(单位时间内处理请求的数目)。
如订单系统,在下单的时候就会往数据库写数据。但是数据库只能支撑每秒1000左右的并发写入,并发量再高就容易宕机。低峰期的时候并发也就100多个,但是在高峰期时候,并发量会突然激增到5000以上,这个时候数据库肯定卡死了。
消息被MQ保存起来了,然后系统就可以按照自己的消费能力来消费,比如每秒1000个消息,这样慢慢写入数据库,这样就不会卡死数据库了。
但是使用了MQ之后,限制消费消息的速度为1000,但是这样一来,高峰期产生的数据势必会被积压在MQ中,高峰就被“削”掉了。但是因为消息积压,在高峰期过后的一段时间内,消费消息的速度还是会维持在1000QPS,直到消费完积压的消息,这就叫做“填谷”
系统可用性降低
系统引入的外部依赖越多,系统稳定性越差。一旦 MQ 宕机,就会对业务造成影响。如何保证MQ的高可用?
系统复杂度提高
MQ 的加入大大增加了系统的复杂度,以前系统间是同步的远程调用,现在是通过 MQ 进行异步调用。如何保证消息没有被重复消费?怎么处理消息丢失情况?那么保证消息传递的顺序性?
一致性问题
A 系统处理完业务,通过 MQ 给B、C、D三个系统发消息,如果 B 系统、C 系统处理成功,D 系统处理失败。如何保证消息数据处理的一致性?
目前业界有很多的 MQ 产品,例如 RabbitMQ、RocketMQ、ActiveMQ、Kafka、ZeroMQ、MetaMq等,也有直接使用 Redis 充当消息队列的案例,而这些消息队列产品,各有侧重,在实际选型时,需要结合自身需求及 MQ 产品特征,综合考虑。
RabbitMQ | ActiveMQ | RocketMQ | Kafka | |
---|---|---|---|---|
公司/ 社区 | Rabbit | Apache | 阿里 | Apache |
开发语言 | Erlang | Java | Java | Scala&Java |
协议支持 | AMQP,XMPP,SMTP,STOMP | OpenWire,STOMP,REST,XMPP,AMQP | 自定义 | 自定义协议,社区封装了http协议支持 |
客户端支持语言 | 官方支持Erlang,Java,Ruby等,社区产出多种API,几乎支持所有语言 | Java,C,C++,Python,PHP,Perl,.net等 | Java,C++(不成熟) | 官方支持Java,社区产出多种API,如PHP,Python等 |
单机吞吐量 | 万级(其次) | 万级(最差) | 十万级(最好) | 十万级(次之) |
消息延迟 | 微妙级 | 毫秒级 | 毫秒级 | 毫秒以内 |
功能特性 | 并发能力强,性能极其好,延时低,社区活跃,管理界面丰富 | 老牌产品,成熟度高,文档较多 | MQ功能比较完备,扩展性佳 | 只支持主要的MQ功能,毕竟是为大数据领域准备的。 |
实现MQ的大致有两种主流方式:AMQP、JMS。
AMQP
AMQP,即 Advanced Message Queuing Protocol(高级消息队列协议),是一个网络协议,是应用层协议的一个开放标准,为面向消息的中间件设计。基于此协议的客户端与消息中间件可传递消息,遵循此协议,不收客户端和中间件产品和开发语言限制。2006年,AMQP 规范发布。类比HTTP。
JMS
JMS 即 Java 消息服务(JavaMessage Service)应用程序接口,是一个 Java 平台中关于面向消息中间件的API
JMS 是 JavaEE 规范中的一种,类比JDBC
很多消息中间件都实现了JMS规范,例如:ActiveMQ。RabbitMQ 官方没有提供 JMS 的实现包,但是开源社区有
AMQP 与 JMS 区别
RabbitMQ官方地址:http://www.rabbitmq.com/
2007年,Rabbit 技术公司基于 AMQP 标准开发的 RabbitMQ 1.0 发布。RabbitMQ 采用 Erlang 语言开
发。Erlang 语言专门为开发高并发和分布式系统的一种语言,在电信领域使用广泛。
RabbitMQ 基础架构如下图:
RabbitMQ 中的相关概念:
RabbitMQ提供了6种模式:简单模式,work模式,Publish/Subscribe发布与订阅模式,Routing路由模式,Topics主题模式,RPC远程调用模式(远程调用,不太算MQ;暂不作介绍);
官网对应模式介绍:https://www.rabbitmq.com/getstarted.html
RabbitMQ在安装号后,可以访问 http://ip:15672
,用默认的guest/guest
用户密码进行登陆,创建用户也可以用管控台添加,如下图:
角色说明:
像mysql拥有数据库的概念并且可以指定⽤户对库和表等操作的权限。RabbitMQ也有类似的权限管理;在RabbitMQ中可以虚拟消息服务器Virtual Host,每个Virtual Hosts相当于⼀个相对独⽴的RabbitMQ服务器,每个VirtualHost之间是相互隔离的。exchange、queue、message不能互通。 相当于mysql的db。Virtual Name⼀般以/开头。
点击Virtual Hosts的名称
展开Permissions,选择设置当前Virtual Hosts的拥有者
简单模式
在上图的模型中,有以下概念:
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.13.1</version>
</dependency>
编写消息生产者com.flaw.rabbitmq.simple.Producer
package com.flaw.rabbitmq.simple; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; public class Producer { static final String QUEUE_NAME = "simple_queue"; public static void main(String[] args) throws Exception { // 创建连接工厂 ConnectionFactory connectionFactory=new ConnectionFactory(); // 主机地址;默认为 localhost connectionFactory.setHost("localhost"); // 连接端口;默认为 5672 connectionFactory.setPort(5672); // 虚拟主机名称;默认为 / connectionFactory.setVirtualHost("/xzk"); // 连接用户名;默认为 guest connectionFactory.setUsername("flaw"); // 连接用户名;默认为 guest connectionFactory.setPassword("flaw"); // 创建连接 Connection connection=connectionFactory.newConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明(创建)队列 /** * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(QUEUE_NAME,true,false,false,null); // 要发送的信息 String message="hello,RabbitMQ!"; /** * exchange:交换机名称,如果没有指定则使用默认Default Exchange * routingKey:路由key,简单模式可以传递队列名称 * props:消息其它属性 * body:消息内容 */ channel.basicPublish("",QUEUE_NAME,null,message.getBytes()); System.out.println("以发送消息:"+message); // 关闭释放资源 channel.close(); connection.close(); } }
在执行上述的消息发送之后;可以登录rabbitMQ的管理控制台,可以发现队列和其消息:
抽取创建connection的工具类com.flaw.rabbitmq.util.ConnectionUtil;
package com.flaw.rabbitmq.util; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; public class ConnectionUtil { public static Connection getConnection() throws Exception { // 创建连接工厂 ConnectionFactory connectionFactory=new ConnectionFactory(); // 主机地址;默认为 localhost connectionFactory.setHost("localhost"); // 连接端口;默认为 5672 connectionFactory.setPort(5672); // 虚拟主机名称;默认为 / connectionFactory.setVirtualHost("/xzk"); // 连接用户名;默认为 guest connectionFactory.setUsername("flaw"); // 连接用户名;默认为 guest connectionFactory.setPassword("flaw"); //创建连接返回 return connectionFactory.newConnection(); } }
编写消息的消费者com.flaw.rabbitmq.simple.Consumer
package com.flaw.rabbitmq.simple; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明(创建)队列 /** * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.QUEUE_NAME,true,false,false,null); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("接收到的消息为:" + new String(body, "utf-8")); } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.QUEUE_NAME, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
上述的入门案例中中其实使用的是如下的简单模式:
在上图的模型中,有以下概念:
Work Queues
与入门程序的 简单模式 相比,多了一个或一些消费端,多个消费端共同消费同一个队列中的消息。
应用场景: 对于 任务过重或任务较多情况使用工作队列可以提高任务处理的速度。
Work Queues
与入门程序的 简单模式 的代码是几乎一样的;可以完全复制,并复制多一个消费者进行多个消费者同时消费消息的测试。
生产者
package com.flaw.rabbitmq.work; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; public class Producer { static final String QUEUE_NAME = "work_queue"; public static void main(String[] args) throws Exception { // 创建连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明(创建)队列 /** * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(QUEUE_NAME,true,false,false,null); for (int i = 1; i <= 30; i++) { // 要发送的信息 String message="hello,RabbitMQ! work模式==="+i; /** * exchange:交换机名称,如果没有指定则使用默认Default Exchange * routingKey:路由key,简单模式可以传递队列名称 * props:消息其它属性 * body:消息内容 */ channel.basicPublish("",QUEUE_NAME,null,message.getBytes()); System.out.println("已发送消息:"+message); } // 关闭释放资源 channel.close(); connection.close(); } }
消费者1
package com.flaw.rabbitmq.work; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer1 { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明(创建)队列 /** * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.QUEUE_NAME,true,false,false,null); // 一次只能接收并处理一个消息 channel.basicQos(1); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { try { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8")); // 休眠1秒,测试时观察更直观 Thread.sleep(1000); //确认消息 channel.basicAck(envelope.getDeliveryTag(), false); } catch (InterruptedException e) { e.printStackTrace(); } } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.QUEUE_NAME, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
消费者2
package com.flaw.rabbitmq.work; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer2 { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明(创建)队列 /** * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.QUEUE_NAME,true,false,false,null); // 一次只能接收并处理一个消息 channel.basicQos(1); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { try { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8")); // 休眠1秒,测试时观察更直观 Thread.sleep(1000); //确认消息 channel.basicAck(envelope.getDeliveryTag(), false); } catch (InterruptedException e) { e.printStackTrace(); } } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.QUEUE_NAME, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
启动两个消费者,然后再启动生产者发送消息;到IDEA的两个消费者对应的控制台查看是否竞争性的接收到消息。
在一个队列中如果有多个消费者,那么消费者之间对于同一个消息的关系是竞争的关系。
订阅模式示例图:
前面2个案例中,只有3个角色:
而在订阅模型中,多了一个exchange角色,而且过程略有变化:
Exchange(交换机)只负责转发消息,不具备存储消息的能力,因此如果没有任何队列与Exchange绑定,或者没有符合路由规则的队列,那么消息会丢失!
发布订阅模式:
生产者:
package com.flaw.rabbitmq.ps; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; /** * 发布与订阅使用的交换机类型为:fanout */ public class Producer { /** * 交换机名称 */ static final String FANOUT_EXCHANGE = "fanout_exchange"; /** * 队列名称 */ static final String FANOUT_QUEUE_1 = "fanout_queue_1"; /** * 队列名称 */ static final String FANOUT_QUEUE_2 = "fanout_queue_2"; public static void main(String[] args) throws Exception { // 创建连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); /** * 声明交换机 * exchange:交换机名称 * type:交换机类型 DIRECT("direct"), FANOUT("fanout"), TOPIC("topic"), HEADERS("headers"); */ channel.exchangeDeclare(FANOUT_EXCHANGE, BuiltinExchangeType.FANOUT); /** * 声明队列 * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(FANOUT_QUEUE_1,true,false,false,null); channel.queueDeclare(FANOUT_QUEUE_2,true,false,false,null); /** * 队列绑定交换机 * queue: 队列名称 * exchange: 交换机名称 * routingKey: 路由键 */ channel.queueBind(FANOUT_QUEUE_1, FANOUT_EXCHANGE, ""); channel.queueBind(FANOUT_QUEUE_2, FANOUT_EXCHANGE, ""); for (int i = 1; i <= 10; i++) { // 要发送的信息 String message="hello,RabbitMQ! 发布订阅模式==="+i; /** * exchange:交换机名称,如果没有指定则使用默认Default Exchange * routingKey:路由key,简单模式可以传递队列名称 * props:消息其它属性 * body:消息内容 */ channel.basicPublish(FANOUT_EXCHANGE,"",null,message.getBytes()); System.out.println("已发送消息:"+message); } // 关闭释放资源 channel.close(); connection.close(); } }
消费者1:
package com.flaw.rabbitmq.ps; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer1 { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明交换机 channel.exchangeDeclare(Producer.FANOUT_EXCHANGE, BuiltinExchangeType.FANOUT); /** * 声明(创建)队列 * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.FANOUT_QUEUE_1,true,false,false,null); // 队列绑定交换机 channel.queueBind(Producer.FANOUT_QUEUE_1, Producer.FANOUT_EXCHANGE, ""); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8")); } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.FANOUT_QUEUE_1, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
消费者2:
package com.flaw.rabbitmq.ps; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer2 { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明交换机 channel.exchangeDeclare(Producer.FANOUT_EXCHANGE, BuiltinExchangeType.FANOUT); /** * 声明(创建)队列 * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.FANOUT_QUEUE_2,true,false,false,null); // 队列绑定交换机 channel.queueBind(Producer.FANOUT_QUEUE_2, Producer.FANOUT_EXCHANGE, ""); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8")); } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.FANOUT_QUEUE_2, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
启动所有消费者,然后使用生产者发送消息;在每个消费者对应的控制台可以查看到生产者发送的所有消息;到达广播(fanout模式)的效果。
在执行完测试代码后,其实到RabbitMQ的管理后台找到 Exchanges 选项卡,点击 fanout_exchange的交换机,可以查看到如下的绑定:
交换机需要与队列进行绑定,绑定之后;一个消息可以被多个消费者都收到。
发布订阅模式与工作队列模式的区别:
路由模式特点:
RoutingKey
(路由key)RoutingKey
。RoutingKey
进行判断,只有队列的Routingkey
与消息的 Routingkey
完全一致,才会接收到消息在编码上与 Publish/Subscribe发布与订阅模式的区别是交换机的类型为:Direct,还有队列绑定交换机的时候需要指定routingkey。
生产者:
package com.flaw.rabbitmq.routing; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; /** * 发布与订阅使用的交换机类型为:direct */ public class Producer { /** * 交换机名称 */ static final String DIRECT_EXCHANGE = "direct_exchange"; /** * 队列名称 */ static final String DIRECT_QUEUE_INSERT = "direct_queue_insert"; /** * 队列名称 */ static final String DIRECT_QUEUE_UPDATE = "direct_queue_update"; public static void main(String[] args) throws Exception { // 创建连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); /** * 声明交换机 * exchange:交换机名称 * type:交换机类型 DIRECT("direct"), FANOUT("fanout"), TOPIC("topic"), HEADERS("headers"); */ channel.exchangeDeclare(DIRECT_EXCHANGE, BuiltinExchangeType.DIRECT); /** * 声明队列 * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(DIRECT_QUEUE_INSERT,true,false,false,null); channel.queueDeclare(DIRECT_QUEUE_UPDATE,true,false,false,null); /** * 队列绑定交换机 * queue: 队列名称 * exchange: 交换机名称 * routingKey: 路由键 */ channel.queueBind(DIRECT_QUEUE_INSERT, DIRECT_EXCHANGE, "insert"); channel.queueBind(DIRECT_QUEUE_UPDATE, DIRECT_EXCHANGE, "update"); // 要发送的信息 String message="新增了商品! 路由模式:routingKey为 insert"; /** * exchange:交换机名称,如果没有指定则使用默认Default Exchange * routingKey:路由key,简单模式可以传递队列名称 * props:消息其它属性 * body:消息内容 */ channel.basicPublish(DIRECT_EXCHANGE,"insert",null,message.getBytes()); System.out.println("已发送消息:"+message); // 要发送的信息 message="修改了商品! 路由模式:routingKey为 update"; /** * exchange:交换机名称,如果没有指定则使用默认Default Exchange * routingKey:路由key,简单模式可以传递队列名称 * props:消息其它属性 * body:消息内容 */ channel.basicPublish(DIRECT_EXCHANGE,"update",null,message.getBytes()); System.out.println("已发送消息:"+message); // 关闭释放资源 channel.close(); connection.close(); } }
消费者1:
package com.flaw.rabbitmq.routing; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer1 { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明交换机 channel.exchangeDeclare(Producer.DIRECT_EXCHANGE, BuiltinExchangeType.DIRECT); /** * 声明(创建)队列 * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.DIRECT_QUEUE_INSERT,true,false,false,null); // 队列绑定交换机 channel.queueBind(Producer.DIRECT_QUEUE_INSERT, Producer.DIRECT_EXCHANGE, "insert"); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8")); } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.DIRECT_QUEUE_INSERT, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
消费者2:
package com.flaw.rabbitmq.routing; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer2 { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明交换机 channel.exchangeDeclare(Producer.DIRECT_EXCHANGE, BuiltinExchangeType.DIRECT); /** * 声明(创建)队列 * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.DIRECT_QUEUE_UPDATE,true,false,false,null); // 队列绑定交换机 channel.queueBind(Producer.DIRECT_QUEUE_UPDATE, Producer.DIRECT_EXCHANGE, "update"); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8")); } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.DIRECT_QUEUE_UPDATE, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
启动所有消费者,然后使用生产者发送消息;在消费者对应的控制台可以查看到生产者发送对应routingKey对应队列的消息;到达按照需要接收的效果。
在执行完测试代码后,其实到RabbitMQ的管理后台找到 Exchanges 选项卡,点击 direct_exchange的交换机,可以查看到如下的绑定:
Routing模式要求队列在绑定交换机时要指定routingKey,消息会转发到符合routingKey的队列。
Topic
类型与 Direct
相比,都是可以根据 RoutingKey
把消息路由到不同的队列。只不过 Topic
类型Exchange 可以让队列在绑定 RoutingKey
的时候使用通配符!
Routingkey
一般都是有一个或多个单词组成,多个单词之间以"."分割,例如: item.insert
通配符规则:
#
:匹配一个或多个词
*
:匹配不多不少恰好1个词
举例:
item.#
:能够匹配 item.insert.abc
或者 item.insert
item.*
:只能匹配 item.insert
图解:
usa.#
,因此凡是以 usa.
开头的 routingKey
都会被匹配到#.news
,因此凡是以 .news
结尾的 routingKey
都会被匹配
生产者:
使用topic类型的Exchange,发送消息的routingKey有3种: goods.insert
、 goods.update
、 goods.delete
:
package com.flaw.rabbitmq.topics; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; /** * 发布与订阅使用的交换机类型为:topic */ public class Producer { /** * 交换机名称 */ static final String TOPIC_EXCHANGE = "topic_exchange"; /** * 队列名称 */ static final String TOPIC_QUEUE_1 = "topic_queue_1"; /** * 队列名称 */ static final String TOPIC_QUEUE_2 = "topic_queue_2"; public static void main(String[] args) throws Exception { // 创建连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); /** * 声明交换机 * exchange:交换机名称 * type:交换机类型 DIRECT("direct"), FANOUT("fanout"), TOPIC("topic"), HEADERS("headers"); */ channel.exchangeDeclare(TOPIC_EXCHANGE, BuiltinExchangeType.TOPIC); // 要发送的信息 String message="新增了商品! Topic模式:routingKey为 goods.insert"; /** * exchange:交换机名称,如果没有指定则使用默认Default Exchange * routingKey:路由key,简单模式可以传递队列名称 * props:消息其它属性 * body:消息内容 */ channel.basicPublish(TOPIC_EXCHANGE,"goods.insert",null,message.getBytes()); System.out.println("已发送消息:"+message); // 要发送的信息 message="修改了商品! Topic模式:routingKey为 goods.update"; /** * exchange:交换机名称,如果没有指定则使用默认Default Exchange * routingKey:路由key,简单模式可以传递队列名称 * props:消息其它属性 * body:消息内容 */ channel.basicPublish(TOPIC_EXCHANGE,"goods.update",null,message.getBytes()); System.out.println("已发送消息:"+message); // 要发送的信息 message="删除了商品! Topic模式:routingKey为 goods.delete"; /** * exchange:交换机名称,如果没有指定则使用默认Default Exchange * routingKey:路由key,简单模式可以传递队列名称 * props:消息其它属性 * body:消息内容 */ channel.basicPublish(TOPIC_EXCHANGE,"goods.delete",null,message.getBytes()); System.out.println("已发送消息:"+message); // 关闭释放资源 channel.close(); connection.close(); } }
消费者1:
接收两种类型的消息:更新商品goods.update
和删除商品goods.delete
package com.flaw.rabbitmq.topics; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer1 { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); //声明交换机 channel.exchangeDeclare(Producer.TOPIC_EXCHANGE, BuiltinExchangeType.TOPIC); /** * 声明(创建)队列 * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.TOPIC_QUEUE_1,true,false,false,null); // 队列绑定交换机 channel.queueBind(Producer.TOPIC_QUEUE_1, Producer.TOPIC_EXCHANGE, "goods.update"); channel.queueBind(Producer.TOPIC_QUEUE_1, Producer.TOPIC_EXCHANGE, "goods.delete"); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8")); } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.TOPIC_QUEUE_1, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
消费者2:
接收所有类型的消息:新增商品,更新商品和删除商品(goods.*
)。
package com.flaw.rabbitmq.topics; import com.flaw.rabbitmq.util.ConnectionUtil; import com.rabbitmq.client.*; import java.io.IOException; public class Consumer2 { public static void main(String[] args) throws Exception { // 获取连接 Connection connection= ConnectionUtil.getConnection(); // 创建频道 Channel channel = connection.createChannel(); // 声明交换机 channel.exchangeDeclare(Producer.TOPIC_EXCHANGE, BuiltinExchangeType.TOPIC); /** * 声明(创建)队列 * queue:队列名称 * durable:是否定义持久化队列 * exclusive:是否独占本次连接 * autoDelete:是否在不使用的时候自动删除队列 * arguments:队列其它参数 */ channel.queueDeclare(Producer.TOPIC_QUEUE_2,true,false,false,null); // 队列绑定交换机 channel.queueBind(Producer.TOPIC_QUEUE_2, Producer.TOPIC_EXCHANGE, "goods.*"); // 创建消费者,并设置消息处理 DefaultConsumer consumer=new DefaultConsumer(channel){ /** * * @param consumerTag 消息者标签,在channel.basicConsume时候可以指定 * @param envelope 消息包的内容,可从中获取消息id,消息routingKey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) * @param properties 属性信息 * @param body 消息体 */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { //路由key System.out.println("路由key为:" + envelope.getRoutingKey()); //交换机 System.out.println("交换机为:" + envelope.getExchange()); //消息id System.out.println("消息id为:" + envelope.getDeliveryTag()); //收到的消息 System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8")); } }; // 监听消息 /** * queue:队列名称 * autoAck:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复 会删除消息,设置为false则需要手动确认 * callback:消息接收到后回调 */ channel.basicConsume(Producer.TOPIC_QUEUE_2, true, consumer); //需要一直监听消息,所以不需要关闭释放资源 } }
启动所有消费者,然后使用生产者发送消息;在消费者对应的控制台可以查看到生产者发送对应routingKey对应队列的消息;到达按照需要接收的效果;并且这些routingKey可以使用通配符。
在执行完测试代码后,其实到RabbitMQ的管理后台找到 Exchanges 选项卡,点击 topic_exchange的交换机,可以查看到如下的绑定:
Topic主题模式可以实现 Publish/Subscribe发布与订阅模式
和 Routing路由模式
的功能;只是Topic在配置routingKey 的时候可以使用通配符,显得更加灵活。
RabbitMQ工作模式:
修改pom.xml文件内容为如下:
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.flaw</groupId> <artifactId>spring-rabbitmq-producer</artifactId> <version>1.0-SNAPSHOT</version> <dependencies> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-context</artifactId> <version>5.1.7.RELEASE</version> </dependency> <dependency> <groupId>org.springframework.amqp</groupId> <artifactId>spring-rabbit</artifactId> <version>2.1.8.RELEASE</version> </dependency> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>4.12</version> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-test</artifactId> <version>5.1.7.RELEASE</version> </dependency> </dependencies> </project>
rabbitmq.host=192.168.56.1
rabbitmq.port=5672
rabbitmq.username=flaw
rabbitmq.password=flaw
rabbitmq.virtual-host=/xzk
<?xml version="1.0" encoding="UTF-8"?> <beans xmlns="http://www.springframework.org/schema/beans" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xmlns:context="http://www.springframework.org/schema/context" xmlns:rabbit="http://www.springframework.org/schema/rabbit" xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd http://www.springframework.org/schema/context https://www.springframework.org/schema/context/spring-context.xsd http://www.springframework.org/schema/rabbit http://www.springframework.org/schema/rabbit/spring-rabbit.xsd"> <!--加载配置文件--> <context:property-placeholder location="classpath:properties/rabbitmq.properties"/> <!-- 定义rabbitmq connectionFactory --> <rabbit:connection-factory id="connectionFactory" host="${rabbitmq.host}" port="${rabbitmq.port}" username="${rabbitmq.username}" password="${rabbitmq.password}" virtual-host="${rabbitmq.virtual-host}"/> <!--定义管理交换机、队列--> <rabbit:admin connection-factory="connectionFactory"/> <!--定义持久化队列,不存在则自动创建;不绑定到交换机则绑定到默认交换机 默认交换机类型为direct,名字为:"",路由键为队列的名称 --> <rabbit:queue id="spring_queue" name="spring_queue" auto-declare="true"/> <!-- ~~~~~~~~~~~~~~~~~~~~~~~~~~~~广播:所有队列都能收到消息 ~~~~~~~~~~~~~~~~~~~~~~~~~~~~ --> <!--定义广播交换机中的持久化队列,不存在则自动创建--> <rabbit:queue id="spring_fanout_queue_1" name="spring_fanout_queue_1" auto-declare="true"/> <!--定义广播交换机中的持久化队列,不存在则自动创建--> <rabbit:queue id="spring_fanout_queue_2" name="spring_fanout_queue_2" auto-declare="true"/> <!--定义广播类型交换机;并绑定上述两个队列--> <rabbit:fanout-exchange id="spring_fanout_exchange" name="spring_fanout_exchange" auto-declare="true"> <rabbit:bindings> <rabbit:binding queue="spring_fanout_queue_1"/> <rabbit:binding queue="spring_fanout_queue_2"/> </rabbit:bindings> </rabbit:fanout-exchange> <!-- ~~~~~~~~~~~~~~~~~~~~~~~~~~~~通配符;*匹配一个单词,#匹配多个单词 ~~~~~~~~~~~~~~~~~~~~~~~~~~~~ --> <!--定义广播交换机中的持久化队列,不存在则自动创建--> <rabbit:queue id="spring_topic_queue_star" name="spring_topic_queue_star" auto-declare="true"/> <!--定义广播交换机中的持久化队列,不存在则自动创建--> <rabbit:queue id="spring_topic_queue_well" name="spring_topic_queue_well" auto-declare="true"/> <!--定义广播交换机中的持久化队列,不存在则自动创建--> <rabbit:queue id="spring_topic_queue_well2" name="spring_topic_queue_well2" auto-declare="true"/> <rabbit:topic-exchange id="spring_topic_exchange" name="spring_topic_exchange" auto-declare="true"> <rabbit:bindings> <rabbit:binding pattern="lxs.*" queue="spring_topic_queue_star"/> <rabbit:binding pattern="lxs.#" queue="spring_topic_queue_well"/> <rabbit:binding pattern="xzk.#" queue="spring_topic_queue_well2"/> </rabbit:bindings> </rabbit:topic-exchange> <!--定义rabbitTemplate对象操作可以在代码中方便发送消息--> <rabbit:template id="rabbitTemplate" connection-factory="connectionFactory"/> </beans>
创建测试文件 spring-rabbitmq-producer\src\test\java\com\flaw\rabbitmq\ProducerTest.java
package com.flaw.rabbitmq; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration(locations = "classpath:spring/spring-rabbitmq.xml") public class ProducerTest { @Autowired private RabbitTemplate rabbitTemplate; /** * 只发队列消息 * 默认交换机类型:direct * 交换机名称为空,路由键为队列的名称 */ @Test public void queueTest(){ // 路由键与队列同名 rabbitTemplate.convertAndSend("spring_queue","发送队列spring_queue的消息。"); } /** * 发送广播 * 交换机类型: fanout * 绑定到该交换机的队列都能收到消息 */ @Test public void fanoutTest(){ /** * exchange:交换机名称 * routingKey:路由键名称(广播设置为空) * object: 需要发送的消息 */ rabbitTemplate.convertAndSend("spring_fanout_exchange","","发送到spring_fanout_exchange交换机的广播消息"); } /** * 通配符 * 交换机类型: topic * 匹配路由键的通配符,*表示一个单词,#表示多个单词 * 绑定到该交换机的匹配队列能够收到对应消息 */ @Test public void topicTest(){ /** * exchange:交换机名称 * routingKey:路由键名称(广播设置为空) * object: 需要发送的消息 */ rabbitTemplate.convertAndSend("spring_topic_exchange","flaw.haha","发送到spring_topic_exchange交换机flaw.haha的消息"); rabbitTemplate.convertAndSend("spring_topic_exchange","flaw.haha.1","发送到spring_topic_exchange交换机flaw.haha.1的消息"); rabbitTemplate.convertAndSend("spring_topic_exchange","flaw.haha.2","发送到spring_topic_exchange交换机flaw.haha.2的消息"); rabbitTemplate.convertAndSend("spring_topic_exchange","xzk.com","发送到spring_topic_exchange交换机xzk.com的消息"); } }
修改pom.xml文件内容为如下:
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.flaw</groupId> <artifactId>spring-rabbitmq-consumer</artifactId> <version>1.0-SNAPSHOT</version> <dependencies> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-context</artifactId> <version>5.1.7.RELEASE</version> </dependency> <dependency> <groupId>org.springframework.amqp</groupId> <artifactId>spring-rabbit</artifactId> <version>2.1.8.RELEASE</version> </dependency> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>4.12</version> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-test</artifactId> <version>5.1.7.RELEASE</version> </dependency> </dependencies> </project>
rabbitmq.host=192.168.56.1
rabbitmq.port=5672
rabbitmq.username=flaw
rabbitmq.password=flaw
rabbitmq.virtual-host=/xzk
在这里插入代码片
spring-rabbitmq-consumer\src\main\java\com\flaw\rabbitmq\listener\SpringQueueListener.java
监听spring_queue队列的消息package com.flaw.rabbitmq.listener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; public class SpringQueueListener implements MessageListener { public void onMessage(Message message) { try { String msg = new String(message.getBody(), "utf-8"); System.out.printf("接收路由名称为:%s,路由键为:%s,队列名为:%s的消息:%s \n", message.getMessageProperties().getReceivedExchange(), message.getMessageProperties().getReceivedRoutingKey(), message.getMessageProperties().getConsumerQueue(), msg); } catch (Exception e) { e.printStackTrace(); } } }
spring-rabbitmq-consumer\src\main\java\com\flaw\rabbitmq\listener\FanoutListener1.java
监听spring_fanout_queue_1队列的消息package com.flaw.rabbitmq.listener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; public class FanoutListener1 implements MessageListener { @Override public void onMessage(Message message) { try { String msg = new String(message.getBody(), "utf-8"); System.out.printf("广播监听器1:接收路由名称为:%s,路由键为:%s,队列名为:%s的消息为:%s \n", message.getMessageProperties().getReceivedExchange(), message.getMessageProperties().getReceivedRoutingKey(), message.getMessageProperties().getConsumerQueue(), msg); } catch (Exception e) { e.printStackTrace(); } } }
spring-rabbitmq-consumer\src\main\java\com\flaw\rabbitmq\listener\FanoutListener2.java
监听spring_fanout_queue_2的消息package com.flaw.rabbitmq.listener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; public class FanoutListener2 implements MessageListener { @Override public void onMessage(Message message) { try { String msg = new String(message.getBody(), "utf-8"); System.out.printf("广播监听器2:接收路由名称为:%s,路由键为:%s,队列名为: %s的消息:%s \n", message.getMessageProperties().getReceivedExchange(), message.getMessageProperties().getReceivedRoutingKey(), message.getMessageProperties().getConsumerQueue(), msg); } catch (Exception e) { e.printStackTrace(); } } }
*
通配符监听器spring-rabbitmq-consumer\src\main\java\com\flaw\rabbitmq\listener\TopicListenerStar.java
监听spring_topic_queue_star队列的消息package com.flaw.rabbitmq.listener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; public class TopicListenerStar implements MessageListener { @Override public void onMessage(Message message) { try { String msg = new String(message.getBody(), "utf-8"); System.out.printf("通配符*监听器:接收路由名称为:%s,路由键为:%s,队列名 为:%s的消息:%s \n", message.getMessageProperties().getReceivedExchange(), message.getMessageProperties().getReceivedRoutingKey(), message.getMessageProperties().getConsumerQueue(), msg); } catch (Exception e) { e.printStackTrace(); } } }
#
通配符监听器1spring-rabbitmq-consumer\src\main\java\com\flaw\rabbitmq\listener\TopicListenerWell1.java
spring_topic_queue_well队列的消息package com.flaw.rabbitmq.listener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; public class TopicListenerWell1 implements MessageListener { @Override public void onMessage(Message message) { try { String msg = new String(message.getBody(), "utf-8"); System.out.printf("通配符#监听器:接收路由名称为:%s,路由键为:%s,队列名 为:%s的消息:%s \n", message.getMessageProperties().getReceivedExchange(), message.getMessageProperties().getReceivedRoutingKey(), message.getMessageProperties().getConsumerQueue(), msg); } catch (Exception e) { e.printStackTrace(); } } }
#
通配符监听器2spring-rabbitmq-consumer\src\main\java\com\flaw\rabbitmq\listener\TopicListenerWell2.java
监听spring_topic_queue_well2队列的消息package com.flaw.rabbitmq.listener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; public class TopicListenerWell2 implements MessageListener { @Override public void onMessage(Message message) { try { String msg = new String(message.getBody(), "utf-8"); System.out.printf("通配符#监听器2:接收路由名称为:%s,路由键为:%s,队列名 为:%s的消息:%s \n", message.getMessageProperties().getReceivedExchange(), message.getMessageProperties().getReceivedRoutingKey(), message.getMessageProperties().getConsumerQueue(), msg); } catch (Exception e) { e.printStackTrace(); } } }
创建测试类:
package com.flaw.rabbitmq; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration(locations = "classpath:spring/spring-rabbitmq.xml") public class ConsumerTest { @Test public void test(){ while (true){ // 死循环让consumer保持连接状态,一直监听队列里面的消息 } } }
在Spring项目中,可以使用Spring-Rabbit去操作RabbitMQ https://github.com/spring-projects/spring-amqp
尤其是在spring boot项目中只需要引入对应的amqp启动器依赖即可,方便的使用RabbitTemplate发送消息,使用注解接收消息。
一般在开发过程中:
生产者工程:
消费者工程:
创建生产者工程springboot-rabbitmq-producer
修改pom.xml文件内容为如下:
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.flaw</groupId> <artifactId>springboot-rabbitmq-producer</artifactId> <version>1.0-SNAPSHOT</version> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.3.12.RELEASE</version> </parent> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> </dependency> </dependencies> </project>
package com.flaw.rabbitmq;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class ProducerApplication {
public static void main(String[] args) {
SpringApplication.run(ProducerApplication.class);
}
}
spring:
rabbitmq:
host: 192.168.56.1
port: 5672
virtual-host: /xzk
username: flaw
password: flaw
package com.flaw.rabbitmq.config; import org.springframework.amqp.core.*; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitmqConfig { /** * 交换机名称 */ public static final String ITEM_TOPIC_EXCHANGE = "springboot_item_topic_exchange"; /** * 队列名称 */ public static final String ITEM_QUEUE = "springboot_item_queue"; /** * 声明交换机 * @return */ @Bean("itemTopicExchange") public Exchange topicExchange(){ return ExchangeBuilder.topicExchange(ITEM_TOPIC_EXCHANGE).durable(true).build(); } /** * 声明队列 * @return */ @Bean("itemQueue") public Queue itemQueue(){ return QueueBuilder.durable(ITEM_QUEUE).build(); } /** * 绑定队列和交换机 * @param queue 队列 * @param exchange 交换机 * @return */ @Bean public Binding itemQueueExchange(@Qualifier("itemQueue") Queue queue, @Qualifier("itemTopicExchange") Exchange exchange){ return BindingBuilder.bind(queue).to(exchange).with("item.#").noargs(); } }
在生产者工程springboot-rabbitmq-producer中创建测试类,发送消息:
package com.flaw.rabbitmq; import com.flaw.rabbitmq.config.RabbitmqConfig; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.test.context.junit4.SpringRunner; @RunWith(SpringRunner.class) @SpringBootTest public class RabbitMQTest { @Autowired private RabbitTemplate rabbitTemplate; @Test public void test(){ rabbitTemplate.convertAndSend(RabbitmqConfig.ITEM_TOPIC_EXCHANGE, "item.insert", "商品新增,routingKey 为item.insert"); rabbitTemplate.convertAndSend(RabbitmqConfig.ITEM_TOPIC_EXCHANGE, "item.update", "商品修改,routingKey 为item.update"); rabbitTemplate.convertAndSend(RabbitmqConfig.ITEM_TOPIC_EXCHANGE, "item.delete", "商品删除,routingKey 为item.delete"); } }
先运行上述测试程序(交换机和队列才能先被声明和绑定),然后启动消费者;在消费者工程springboot-rabbitmq-consumer中控制台查看是否接收到对应消息。
另外,也可以在RabbitMQ的管理控制台中查看到交换机与队列的绑定:
创建消费者工程springboot-rabbitmq-consumer
修改pom.xml文件内容为如下:
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.flaw</groupId> <artifactId>springboot-rabbitmq-consumer</artifactId> <version>1.0-SNAPSHOT</version> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.3.12.RELEASE</version> </parent> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> </dependency> </dependencies> </project>
package com.flaw.rabbitmq;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class ConsumerApplication {
public static void main(String[] args) {
SpringApplication.run(ConsumerApplication.class);
}
}
创建application.yml,内容如下:
spring:
rabbitmq:
host: 192.168.56.1
port: 5672
virtual-host: /xzk
username: flaw
password: flaw
编写消息监听器com.flaw.rabbitmq.listener.MyListener
package com.flaw.rabbitmq.listener; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; @Component public class MyListener { /** * 监听某个队列的消息 * @param message 接收的消息 */ @RabbitListener(queues = "springboot_item_queue") public void myListener1(String message){ System.out.println("消费者收到的消息:"+message); } }
package com.flaw.rabbitmq; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.test.context.junit4.SpringRunner; /** * @author Administrator * @since 2021/10/27 10:41 */ @RunWith(SpringRunner.class) @SpringBootTest public class ConsumerTest { @Test public void test(){ while (true){ // 死循环让consumer保持连接状态,一直监听队列里面的消息 } } }
在使用 RabbitMQ 的时候,作为消息发送方希望杜绝任何消息丢失或者投递失败场景。RabbitMQ 为我们提供了两种方式用来控制消息的投递可靠性模式。
rabbitmq 整个消息投递的路径为:
producer—>rabbitmq broker—>exchange—>queue—>consumer
我们将利用这两个 callback 控制消息的可靠性投递
创建producer项目
添加依赖
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.flaw</groupId> <artifactId>spring-rabbitmq-senior-producer</artifactId> <version>1.0-SNAPSHOT</version> <dependencies> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-context</artifactId> <version>5.1.7.RELEASE</version> </dependency> <dependency> <groupId>org.springframework.amqp</groupId> <artifactId>spring-rabbit</artifactId> <version>2.1.8.RELEASE</version> </dependency> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>4.12</version> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-test</artifactId> <version>5.1.7.RELEASE</version> </dependency> </dependencies> </project>
创建rabbitmq.properties配置文件
rabbitmq.host=127.0.0.1
rabbitmq.port=5672
rabbitmq.username=guest
rabbitmq.password=guest
rabbitmq.virtual-host=/
创建spring-rabbitmq-producer.xml配置文件
<?xml version="1.0" encoding="UTF-8"?> <beans xmlns="http://www.springframework.org/schema/beans" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xmlns:context="http://www.springframework.org/schema/context" xmlns:rabbit="http://www.springframework.org/schema/rabbit" xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd http://www.springframework.org/schema/context https://www.springframework.org/schema/context/spring-context.xsd http://www.springframework.org/schema/rabbit http://www.springframework.org/schema/rabbit/spring-rabbit.xsd"> <!--加载配置文件--> <context:property-placeholder location="classpath*:rabbitmq.properties" /> <!--定义rabbitmq连接工厂connectionFactory--> <!-- publisher-confirms="true" 开启确认模式,默认为false --> <rabbit:connection-factory id="connectionFactory" host="${rabbitmq.host}" port="${rabbitmq.port}" virtual-host="${rabbitmq.virtual-host}" username="${rabbitmq.username}" password="${rabbitmq.password}" publisher-confirms="true" /> <!--定义管理交换机、队列--> <rabbit:admin connection-factory="connectionFactory"/> <!--定义rabbitTemplate对象操作可以在代码中方便发送消息--> <rabbit:template id="rabbitTemplate" connection-factory="connectionFactory"/> <!--消息可靠性投递(生产端producer)--> <rabbit:queue id="test_queue_confirm" name="test_queue_confirm"></rabbit:queue> <rabbit:direct-exchange name="test_exchange_confirm"> <rabbit:bindings> <rabbit:binding queue="test_queue_confirm" key="confirm"></rabbit:binding> </rabbit:bindings> </rabbit:direct-exchange> </beans>
消息从 producer 到 exchange 则会返回一个 confirmCallback
package com.flaw.rabbitmq; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.amqp.rabbit.connection.CorrelationData; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration(locations = "classpath:spring-rabbitmq-producer.xml") public class ProducerTest { @Autowired private RabbitTemplate rabbitTemplate; /** * 确认模式 * 步骤: * 1. 确认模式开启:在spring-rabbitmq-producer.xml connectionFactory中开启publisher-confirms="true" * 2. 在rabbitTemplate定义ConfirmCallBack回调函数 * 3. 发送消息 */ @Test public void testConfirm() { // 定义回调函数 rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() { /** * * @param correlationData 回调的相关数据 * @param ack exchange交换机是否成功收到了消息。true(ack) 成功,false(nack)代表失败 * @param cause 失败原因 */ @Override public void confirm(CorrelationData correlationData, boolean ack, String cause) { System.out.println("confirm方法被执行了。。。"); if (ack) { System.out.println("接收消息成功"); }else { System.out.println("接收消息失败,失败原因:"+cause); } } }); // 发送消息 //rabbitTemplate.convertAndSend("test_exchange_confirm", "confirm", "the message a confirm message...."); // 失败消息发送 rabbitTemplate.convertAndSend("test_exchange_confirm111", "confirm", "the message a confirm message...."); } }
消息从 exchange–>queue 投递失败则会返回一个 returnCallback
/** * 回退模式: 当消息发送给Exchange后,Exchange路由到Queue失败时才会执行 ReturnCallBack * 步骤: * 1. 开启回退模式: 在spring-rabbitmq-producer.xml connectionFactory中开启publisher-returns="true" * 2. 定义回调 * 3. 设置Exchange处理消息失败的模式:setMandatory * * * 如果消息没有路由到Queue,则丢弃消息(默认) * 如果消息没有路由到Queue,返回给消息发送方ReturnCallBack */ @Test public void testReturn(){ // 设置交换机处理失败消息的模式 rabbitTemplate.setMandatory(true); // 定义回调函数 rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() { /** * * @param message 返回消息对象 * @param replyCode 错误码 * @param replyText 错误信息 * @param exchange 交换机 * @param routingKey 路由键 */ @Override public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) { System.out.println("returnedMessage方法执行了。。。"); System.out.println(message); System.out.println(replyCode); System.out.println(replyText); System.out.println(exchange); System.out.println(routingKey); } }); // 发送信息 rabbitTemplate.convertAndSend("test_exchange_confirm", "confirm111", "the message a returnCallback message...."); }
ack指Acknowledge,确认。表示消费端收到消息后的确认方式。
有三种确认方式:
其中自动确认是指,当消息一旦被Consumer接收到,则自动确认收到,并将相应 message 从RabbitMQ 的消息缓存中移除。但是在实际业务处理中,很可能消息接收到,业务处理出现异常,那么该消息就会丢失。如果设置了手动确认方式,则需要在业务处理成功后,调用channel.basicAck(),手动签收,如果出现异常,则调用channel.basicNack()方法,让其自动重新发送消息。
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.flaw</groupId> <artifactId>spring-rabbitmq-senior-consumer</artifactId> <version>1.0-SNAPSHOT</version> <dependencies> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-context</artifactId> <version>5.1.7.RELEASE</version> </dependency> <dependency> <groupId>org.springframework.amqp</groupId> <artifactId>spring-rabbit</artifactId> <version>2.1.8.RELEASE</version> </dependency> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>4.12</version> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-test</artifactId> <version>5.1.7.RELEASE</version> </dependency> </dependencies> </project>
创建 spring-rabbitmq-senior-consumer\src\main\resources\rabbitmq.properties
配置文件
rabbitmq.host=127.0.0.1
rabbitmq.port=5672
rabbitmq.username=guest
rabbitmq.password=guest
rabbitmq.virtual-host=/
创建 spring-rabbitmq-senior-consumer\src\main\resources\spring-rabbitmq-consumer.xml
配置文件
<?xml version="1.0" encoding="UTF-8"?> <beans xmlns="http://www.springframework.org/schema/beans" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xmlns:context="http://www.springframework.org/schema/context" xmlns:rabbit="http://www.springframework.org/schema/rabbit" xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd http://www.springframework.org/schema/context https://www.springframework.org/schema/context/spring-context.xsd http://www.springframework.org/schema/rabbit http://www.springframework.org/schema/rabbit/spring-rabbit.xsd"> <!--加载配置文件--> <context:property-placeholder location="classpath:rabbitmq.properties"/> <!-- 定义rabbitmq连接工厂connectionFactory --> <rabbit:connection-factory id="connectionFactory" host="${rabbitmq.host}" port="${rabbitmq.port}" virtual-host="${rabbitmq.virtual-host}" username="${rabbitmq.username}" password="${rabbitmq.password}" /> <!--扫描监听类--> <context:component-scan base-package="com.flaw.rabbitmq.listener" /> <rabbit:listener-container connection-factory="connectionFactory" acknowledge="manual" prefetch="1"> <rabbit:listener ref="ackListener" queue-names="test_queue_confirm"/> </rabbit:listener-container> </beans>
创建 spring-rabbitmq-senior-consumer\src\main\java\com\flaw\rabbitmq\listener\AckListener.java
监听器
package com.flaw.rabbitmq.listener; import com.rabbitmq.client.Channel; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener; import org.springframework.stereotype.Component; /** * Consumer ACK机制: * 1. 设置手动签收。在配置文件中rabbit:listener-container设置 acknowledge="manual" * 2. 让监听器类实现ChannelAwareMessageListener接口 * 3. 如果消息成功处理,则调用channel的 basicAck()签收 * 4. 如果消息处理失败,则调用channel的basicNack()拒绝签收,broker重新发送给consumer * */ @Component public class AckListener implements ChannelAwareMessageListener { @Override public void onMessage(Message message, Channel channel) throws Exception { // 获取当前消息的标签 long deliveryTag = message.getMessageProperties().getDeliveryTag(); try { System.out.println(new String(message.getBody())); System.out.println("处理业务逻辑"); //模拟业务处理异常 int i = 3/0; // 手动签收 channel.basicAck(deliveryTag, true); } catch (Exception e) { //e.printStackTrace(); /** * 拒绝签收 * * deliveryTag: 当前信息的标签 * multiple:true表示拒收所有信息,false表示拒收当前标签的信息 * requeue:true表示拒收的消息重回队列,false表示丢弃拒收的消息或者让拒收的消息成为死信 */ channel.basicNack(deliveryTag,true,true); //channel.basicReject(deliveryTag,true); } } }
创建测试类 spring-rabbitmq-senior-consumer\src\test\java\com\flaw\rabbitmq\ConsumerTest.java
package com.flaw.rabbitmq; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration(locations = "classpath:spring-rabbitmq-consumer.xml") public class ConsumerTest { @Test public void test(){ while (true){ // 死循环让consumer保持连接状态,一直监听队列里面的消息 } } }
当业务出现异常时拒绝签收当前消息,让消息重回队列进行重复执行,等待回复正常后签收。(场景:在网络波动时出现异常就拒收当前消息,让其重回队列重复获取当前消息,等网络回复正常后在签收当前消息)
当消费端有大量请求时(例如秒杀活动),大量的请求写入数据库时会造成数据库宕机。这时可以使用MQ存储消费端请求,让系统从MQ中拉取自身每秒最大能处理的请求量写入数据库,这样就能保证了数据库不会宕机
package com.flaw.rabbitmq.listener; import com.rabbitmq.client.Channel; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener; import org.springframework.stereotype.Component; /** * Consumer 限流机制 * 1. 确保ack机制为手动签收 * 2. 在spring-rabbitmq-consumer.xml中listener-container配置属性prefetch = 1, * 表示消费端每次从mq拉去一条消息来消费,直到手动确认消费完毕后,才会继续拉去下一条消息。 * */ @Component public class QosListener implements ChannelAwareMessageListener { @Override public void onMessage(Message message, Channel channel) throws Exception { // 休眠一秒更能直观感受 Thread.sleep(1000); // 获取消息 System.out.println(new String(message.getBody())); // 处理业务逻辑 // 签收 channel.basicAck(message.getMessageProperties().getDeliveryTag(),true); } }
spring-rabbitmq-consumer.xml
配置
<rabbit:listener-container connection-factory="connectionFactory" acknowledge="manual" prefetch="1">
<rabbit:listener ref="qosListener" queue-names="test_queue_confirm"/>
</rabbit:listener-container>
Time To Live,消息过期时间设置
管控台中设置队列TTL
代码实现(在producer项目中)
在写代码实现时先删除在管控台手动创建的队列,这样好验证代码是否创建成功
配置文件(spring-rabbitmq-producer.xml)
<!--ttl(消息过期时间设置)-->
<rabbit:queue name="test_queue_ttl" id="test_queue_ttl">
<rabbit:queue-arguments>
<entry key="x-message-ttl" value="100000" value-type="java.lang.Integer"></entry>
</rabbit:queue-arguments>
</rabbit:queue>
<rabbit:topic-exchange name="test_exchange_ttl">
<rabbit:bindings>
<rabbit:binding pattern="ttl.#" queue="test_queue_ttl"></rabbit:binding>
</rabbit:bindings>
</rabbit:topic-exchange>
测试代码
/** * ttl:消息过期时间设置 * 1. 队列统一过期 * 2. 消息单独过期 * * 如果设置了消息的过期时间,也设置了队列的过期时间,它以时间短的为准。 * 队列过期后,会将队列所有消息全部移除。 * 消息过期后,只有消息在队列顶端,才会判断其是否过期(移除掉) */ @Test public void testTtl(){ /*for (int i = 1; i <= 10; i++) { // 发送消息 rabbitTemplate.convertAndSend("test_exchange_ttl", "ttl.haha", "message ttl...."+i); }*/ MessagePostProcessor messagePostProcessor = new MessagePostProcessor() { @Override public Message postProcessMessage(Message message) { // 设置消息过期时间5秒 message.getMessageProperties().setExpiration("5000"); // 返回消息 return message; } }; // 消息单独过期 rabbitTemplate.convertAndSend("test_exchange_ttl", "ttl.aaa", "message ttl bbb....", messagePostProcessor); // 注意在队列中设置过期消息不生效,消息过期后,只有消息在队列顶端,才会判断其是否过期(移除掉)。可以运行下面的代码进行验证 /*for (int i = 0; i < 10; i++) { if (i==5){ // 消息单独过期5秒 rabbitTemplate.convertAndSend("test_exchange_ttl", "ttl.bbb", "message ttl bbb....", messagePostProcessor); }else { // 配置文件中的过期时间10秒 rabbitTemplate.convertAndSend("test_exchange_ttl", "ttl.haha", "message ttl...."+i); } }*/ }
死信队列,英文缩写:DLX 。Dead Letter Exchange(死信交换机),当消息成为Dead message后,可以被重新发送到另一个交换机,这个交换机就是DLX。
消息成为死信的三种情况:
队列绑定死信交换机:
给队列设置参数: x-dead-letter-exchange 和 x-dead-letter-routing-key
代码实现:
spring-rabbitmq-producer.xml
配置
<!-- 1. 声明正常的队列(test_queue_dlx)和交换机(test_exchange_dlx) --> <rabbit:queue name="test_queue_dlx" id="test_queue_dlx"> <!--正常队列绑定死信交换机--> <rabbit:queue-arguments> <!--x-dead-letter-exchange 死信交换机名称--> <entry key="x-dead-letter-exchange" value="exchange_dlx" /> <!--x-dead-letter-routing-key 发送给死信交换机的路由key--> <entry key="x-dead-letter-routing-key" value="dlx.haha" /> <!--设置队列的长度限制 max-length--> <entry key="x-max-length" value="10" value-type="java.lang.Integer" /> </rabbit:queue-arguments> </rabbit:queue> <rabbit:topic-exchange name="test_exchange_dlx"> <rabbit:bindings> <rabbit:binding pattern="test.dlx.#" queue="test_queue_dlx" /> </rabbit:bindings> </rabbit:topic-exchange> <!-- 2. 声明死信队列(queue_dlx)和死信交换机(exchange_dlx) --> <rabbit:queue name="queue_dlx" id="queue_dlx" /> <rabbit:topic-exchange name="exchange_dlx"> <rabbit:bindings> <rabbit:binding pattern="dlx.#" queue="queue_dlx" /> </rabbit:bindings> </rabbit:topic-exchange>
测试代码:
生产者测试( ProducerTest.java
测试类)
/** * 发送测试死信消息 * 1. 过期时间 * 2. 长度限制 * 3. 消息拒收 */ @Test public void testDlx(){ // 测试过期时间,死信消息 /*rabbitTemplate.convertAndSend("test_exchange_dlx","test.dlx.haha","我是一条ttl消息,我会死吗?",message -> { // 设置消息过期时间10秒 message.getMessageProperties().setExpiration("10000"); // 返回消息 return message; });*/ // 测试超过长度限制后,消息死信 /*for (int i = 1; i <= 20; i++) { rabbitTemplate.convertAndSend("test_exchange_dlx","test.dlx.haha","我是一条max-length消息,我会死吗?"+i); }*/ // 测试消息拒收(配合消费端ConsumerTest.java测试类测试) rabbitTemplate.convertAndSend("test_exchange_dlx","test.dlx.haha","我是一条nack消息,我会死吗?"); }
消费者监听器 DlxListener.java
package com.flaw.rabbitmq.listener; import com.rabbitmq.client.Channel; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener; import org.springframework.stereotype.Component; /** * Consumer 死信队列 * */ @Component public class DlxListener implements ChannelAwareMessageListener { @Override public void onMessage(Message message, Channel channel) throws Exception { long deliveryTag = message.getMessageProperties().getDeliveryTag(); try { // 接收转换消息 System.out.println(new String(message.getBody())); // 处理业务逻辑 System.out.println("处理业务逻辑..."); // 模拟业务出现异常 int i = 3/0; // 手动签收 channel.basicAck(deliveryTag,true); } catch (Exception e) { //e.printStackTrace(); System.out.println("出现异常,拒绝接受"); // 拒绝签收,不重回队列 requeue=false channel.basicNack(deliveryTag,true,false); } } }
spring-rabbitmq-consumer.xml
配置
<rabbit:listener-container connection-factory="connectionFactory" acknowledge="manual" prefetch="1">
<rabbit:listener ref="dlxListener" queue-names="test_queue_dlx"/>
</rabbit:listener-container>
延迟队列,即消息进入队列后不会立即被消费,只有到达指定时间后,才会被消费。
需求:
实现方式:
很可惜,在RabbitMQ中并未提供延迟队列功能。但是可以使用:TTL+死信队列 组合实现延迟队列的效果。
代码实现:
生产者测试( ProducerTest.java
测试类)
@Test
public void testDelay() throws Exception {
SimpleDateFormat sdf=new SimpleDateFormat("yyyy年MM月dd日HH:mm:ss SSS");
// 发送订单消息。 模拟在订单系统中,下单成功后,10秒内未支付该订单就会进入死信队列
rabbitTemplate.convertAndSend("order_exchange","order.msg","订单信息: id=1,time="+sdf.format(new Date()));
// 打印10秒倒计时 (10秒后订单未处理就会自动发送到死信队列,当启动Consumer测试时,发送的订单信息会被消费不会进入死信队列)
for (int i = 10; i > 0; i--) {
System.out.println(i+"...");
Thread.sleep(1000);
}
}
消费者监听器(OrderListener.java
)
package com.flaw.rabbitmq.listener; import com.rabbitmq.client.Channel; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener; import org.springframework.stereotype.Component; @Component public class OrderListener implements ChannelAwareMessageListener { @Override public void onMessage(Message message, Channel channel) throws Exception { long deliveryTag = message.getMessageProperties().getDeliveryTag(); try { // 接收转换消息 System.out.println(new String(message.getBody())); // 处理业务逻辑 System.out.println("处理业务逻辑..."); System.out.println("订单支付完成..."); // 手动签收 channel.basicAck(deliveryTag,true); } catch (Exception e) { //e.printStackTrace(); System.out.println("出现异常,拒绝接受"); // 拒绝签收,不重回队列 requeue=false channel.basicNack(deliveryTag,true,false); } } }
spring-rabbitmq-consumer.xml
配置
<rabbit:listener-container connection-factory="connectionFactory" acknowledge="manual" prefetch="1">
<rabbit:listener ref="orderListener" queue-names="order_queue"/>
</rabbit:listener-container>
当消费者没启动时,生产者测试产生的订单会在10秒后进入死信队列
启动消费者模拟完成支付业务,生产者产生的订单就会被消费
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。