1.MQ
MQ全称 Message Queue(消息队列),是在消息的传输过程中 保存消息的容器。它是应用程序和应用程序之间的通信方法
1.1 为什么使用MQ
在项目中,可将一些无需即时返回且耗时的操作提取出来,进行异步处理,而这种异步处理的方式大大的节省了服务器的请求响应时间,从而提高了系统的吞吐量。
1.2MQ的好处
1.应用解耦 系统间通过消息通信,不用关心其他系统的处理。
2.异步提速 相比于传统的串行、并行方式,提高了系统吞吐量。
3.削峰填谷 可以通过消息队列长度控制请求量;可以缓解短时间内的高并发请求。
简单来说: 就是在访问量剧增的情况下,但是应用仍然不能停,比如“双十一”下单的人多,但是淘宝这个应用仍然要运行,所以就可以使用消息中间件采用队列的形式减少突然访问的压力
使用MQ后,可以提高系统稳定性
1.3劣势
-
系统可用性降低 系统引入的外部依赖越多,系统稳定性越差。一旦 MQ 宕机,就会对业务造成影响。如何保证MQ的高可用?
-
系统复杂度提高 MQ 的加入大大增加了系统的复杂度,以前系统间是同步的远程调用,现在是通过 MQ 进行异步调用。如何保证消息没有被重复消费?怎么处理消息丢失情况?那么保证消息传递的顺序性?
-
一致性问题 A 系统处理完业务,通过 MQ 给B、C、D三个系统发消息数据,如果 B 系统、C 系统处理成功,D 系统处理失败。如何保证消息数据处理的一致性?
1.4常见的MQ组件
RabbitMQ、RocketMQ、ActiveMQ、Kafka、ZeroMQ、MetaMq等
2.RabbitMQ
RabbitMQ是一个由erlang开发的AMQP(Advanced Message Queue 高级消息队列协议 )的开源实现,由于erlang 语言的高并发特性,性能较好,本质是个队列,FIFO 先入先出,里面存放的内容是message
RabbitMQ是一个消息中间件:它接受并转发消息。你可以把它当做一个快递站点,当你要发送一个包裹时,你把你的包裹放到快递站,快递员最终会把你的快递送到收件人那里,按照这种逻辑RabbitMQ是一个快递站,一个快递员帮你传递快件。RabbitMQ与快递站的主要区别在于,它不处理快件而是接收,存储和转发消息数据。
2.1RabbitMQ的原理
核心组件
-
生产者(Producer):负责发送消息到交换器的客户端应用程序。
-
消费者(Consumer):从队列中获取并处理消息的客户端应用程序。
-
交换器(Exchange):接收生产者发送的消息,并根据路由规则将消息转发到相应的队列。
-
队列(Queue):存储消息,直到消费者取走消息。
-
绑定(Binding):定义交换器和队列之间的关联关系。
工作流程
-
消息发送:生产者通过信道(Channel)将消息发送到交换器。
-
消息路由:交换器根据路由键(Routing Key)和绑定键(Binding Key)将消息路由到相应的队列。
-
消息存储:队列存储消息,等待消费者取走。
-
消息消费:消费者通过信道从队列中获取消息并处理。
交换器类型
-
Direct:根据完全匹配的路由键将消息发送到相应的队列。
-
Fanout:将消息广播到所有绑定的队列,不考虑路由键。
-
Topic:根据模式匹配的路由键将消息发送到相应的队列。
2.2简单模式simple
生产者向队列投递消息,消费者从其中取出消息
1.依赖
<!-- java连接rabbitmq的依赖--><dependency><groupId>com.rabbitmq</groupId><artifactId>amqp-client</artifactId><version>5.16.0</version></dependency>
2.生产消息
package com.ghx.hello;import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;import java.io.IOException;
import java.util.concurrent.TimeoutException;/*** @author :guo* @date :Created in 2025/3/20 11:35* @description:* @version:*/
public class Test01 {public static void main(String[] args) throws IOException, TimeoutException {//创建连接工厂ConnectionFactory factory=new ConnectionFactory();//rabbitmq服务器地址 默认本地localhostfactory.setHost("xxxx");//端口号 默认5672factory.setPort(5672);//用户名 密码 默认guestfactory.setUsername("guest");factory.setPassword("guest");factory.setVirtualHost("/");//创建连接对象Connection connection=factory.newConnection();//获取channel对象Channel channel = connection.createChannel();//创建队列 存在则不创建,不存在则创建//String queue, 队列名// boolean durable, 是否持久化// boolean exclusive, 是否独占队列 false// boolean autoDelete,是否自动删除 false// Map<String, Object> arguments 队列的参数配置--消息的格式 消息存放的时间等channel.queueDeclare("hello",true,false,false,null);String msg="hello rabbitmq2";//String exchange,交换机的名称 "":默认交换机// String routingKey, 路由key "hello":队列名// BasicProperties props, 消息的属性--设置过期时间 设置id等 null// byte[] body 消息的内容channel.basicPublish("","hello",null,msg.getBytes());System.out.println("消息发送成功");channel.close();connection.close();}
}
3.消费消息
package com.ghx.hello;import com.rabbitmq.client.*;import java.io.IOException;
import java.util.concurrent.TimeoutException;/*** @author :guo* @date :Created in 2025/3/20 14:22* @description:* @version:*/
public class Test01 {public static void main(String[] args) throws IOException, TimeoutException {//创建连接工厂ConnectionFactory factory=new ConnectionFactory();//rabbitmq服务器地址 默认本地localhostfactory.setHost("xxxx");//端口号 默认5672factory.setPort(5672);//用户名 密码 默认guestfactory.setUsername("guest");factory.setPassword("guest");//虚拟机名称 默认/factory.setVirtualHost("/");//创建连接对象Connection connection = factory.newConnection();//获取channel对象Channel channel = connection.createChannel();DefaultConsumer consumer=new DefaultConsumer(channel){@Overridepublic void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {System.out.println("消费者1收到消息"+new String(body));}};//接受消息channel.basicConsume("hello",true,consumer);//不要关闭连接和channel 监听消息}
}
2.3工作者模式work queues
多个消费者消费同一个队列中的消息,多个消费者之间属于竞争关系,一个消息只能被一个消费者消费,适合对于任务过重或任务较多的情况,使用工作队列可以提高任务的处理速度
1.生产者
package com.ghx.work;import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;/*** @author :guo* @date :Created in 2025/3/20 14:51* @description:* @version:*/
public class Test03 {private static final String QUEUE_NAME="queue01";public static void main(String[] args) {ConnectionFactory factory=new ConnectionFactory();factory.setHost("xxxx");factory.setPort(5672);factory.setUsername("guest");factory.setPassword("guest");factory.setVirtualHost("/");try {Connection connection = factory.newConnection();Channel channel = connection.createChannel();channel.queueDeclare(QUEUE_NAME,true,false,false,null);for (int i = 0; i < 10; i++){String msg="你好 世界"+i;channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,msg.getBytes("utf-8"));}channel.close();connection.close();}catch (Exception e){}}
}
2. 2个消费者
package com.ghx.work;import com.rabbitmq.client.*;import java.io.IOException;
import java.util.concurrent.TimeoutException;/*** @author :guo* @date :Created in 2025/3/20 15:00* @description:* @version:*/
public class Test03 {private static final String QUEUE_NAME="queue01";public static void main(String[] args) throws IOException, TimeoutException {ConnectionFactory factory=new ConnectionFactory();factory.setHost("xxxX");factory.setPort(5672);factory.setUsername("guest");factory.setPassword("guest");factory.setVirtualHost("/");Connection connection = factory.newConnection();Channel channel = connection.createChannel();channel.queueDeclare(QUEUE_NAME,true,false,false,null);channel.basicQos(1);Consumer consumer = new DefaultConsumer(channel){@Overridepublic void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {System.out.println("消费者1收到消息"+new String(body));}};//接收消息channel.basicConsume(QUEUE_NAME,true,consumer);}}
package com.ghx.work;import com.rabbitmq.client.*;import java.io.IOException;
import java.util.concurrent.TimeoutException;/*** @author :guo* @date :Created in 2025/3/20 15:00* @description:* @version:*/
public class Consumer02 {private static final String QUEUE_NAME="queue01";public static void main(String[] args) throws IOException, TimeoutException {ConnectionFactory factory=new ConnectionFactory();factory.setHost("xxxx");factory.setPort(5672);factory.setUsername("guest");factory.setPassword("guest");factory.setVirtualHost("/");Connection connection = factory.newConnection();Channel channel = connection.createChannel();channel.queueDeclare(QUEUE_NAME,true,false,false,null);channel.basicQos(1);Consumer consumer = new DefaultConsumer(channel){@Overridepublic void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {System.out.println("消费者2收到消息"+new String(body));}};//接收消息channel.basicConsume(QUEUE_NAME,true,consumer);}}
2.3发布订阅模式 publish/subscribe
x : 交换机
一方面,接收生产者发送的消息。另一方面,知道如何处理消息,例如递交给某个特别队列、递交给所有队列、或是将消息丢弃。到底如何操作,取决于Exchange的类型。Exchange有常见以下3种类型:
-
Fanout:广播,将消息交给所有绑定到交换机的队列
-
Direct:定向,把消息交给符合指定routing key 的队列
-
Topic:通配符,把消息交给符合routing pattern(路由模式) 的队列
每个消费者都有自己独立的队列
2.3.1生产者
package com.ghx.work;import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;import java.io.IOException;
import java.util.concurrent.TimeoutException;/*** @author :guo* @date :Created in 2025/3/20 11:35* @description:* @version:*/
public class Test01 {public static void main(String[] args) throws IOException, TimeoutException {//创建连接工厂ConnectionFactory factory=new ConnectionFactory();//rabbitmq服务器地址 默认本地localhostfactory.setHost("xxxx");//端口号 默认5672factory.setPort(5672);//用户名 密码 默认guestfactory.setUsername("guest");factory.setPassword("guest");factory.setVirtualHost("/");//创建连接对象Connection connection=factory.newConnection();//获取channel对象Channel channel = connection.createChannel();//创建交换机
// String exchange,交换机的名称
// BuiltinExchangeType type, 交换机的类型
// boolean durable: 是否持久化channel.exchangeDeclare("fanout_exchange", BuiltinExchangeType.FANOUT,true);//创建队列channel.queueDeclare("fanout_queue1",true,false,false,null);channel.queueDeclare("fanout_queue2",true,false,false,null);//绑定队列和交换机
// String queue,队列名
// String exchange,交换机名
// String routingKey: 路由key 因为广播模式没有路由key ""channel.queueBind("fanout_queue1","fanout_exchange","");channel.queueBind("fanout_queue2","fanout_exchange","");//发送消息String msg="hello fanout交换机";channel.basicPublish("fanout_exchange","",null,msg.getBytes());channel.close();connection.close();}
}
2.4路由模式routing
-
队列与交换机的绑定,不能是任意绑定了,而是要指定一个 RoutingKey(路由key)
-
消息的发送方在向 Exchange 发送消息时,也必须指定消息的 RoutingKey
-
Exchange 不再把消息交给每一个绑定的队列,而是根据消息的 Routing Key 进行判断,只有队列的Routingkey 与消息的 Routing key 完全一致,才会接收到消息
package com.ghx.router;import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;import java.io.IOException;
import java.util.concurrent.TimeoutException;/*** @author :guo* @date :Created in 2025/3/20 11:35* @description:* @version:*/
public class Test01 {public static void main(String[] args) throws IOException, TimeoutException {//创建连接工厂ConnectionFactory factory=new ConnectionFactory();//rabbitmq服务器地址 默认本地localhostfactory.setHost("xxxx");//端口号 默认5672factory.setPort(5672);//用户名 密码 默认guestfactory.setUsername("guest");factory.setPassword("guest");factory.setVirtualHost("/");//创建连接对象Connection connection=factory.newConnection();//获取channel对象Channel channel = connection.createChannel();//创建交换机
// String exchange,交换机的名称
// BuiltinExchangeType type, 交换机的类型
// boolean durable: 是否持久化channel.exchangeDeclare("direct_exchange", BuiltinExchangeType.DIRECT,true);//创建队列channel.queueDeclare("direct_queue1",true,false,false,null);channel.queueDeclare("direct_queue2",true,false,false,null);//绑定队列和交换机
// String queue,队列名
// String exchange,交换机名
// String routingKey: 路由key 因为广播模式没有路由key ""channel.queueBind("direct_queue1","direct_exchange","error");channel.queueBind("direct_queue2","direct_exchange","error");channel.queueBind("direct_queue2","direct_exchange","info");channel.queueBind("direct_queue2","direct_exchange","warning");//发送消息String msg="hello direct交换机";channel.basicPublish("direct_exchange","info",null,msg.getBytes());channel.close();connection.close();}
}
2.5主题模式topics
-
Topic 类型与 Direct 相比,都是可以根据 RoutingKey 把消息路由到不同的队列。只不过 Topic 类型Exchange 可以让队列在绑定 Routing key 的时候使用==通配符==!
-
Routingkey 一般都是有一个或多个单词组成,多个单词之间以”.”分割,例如: item.insert
-
通配符规则:# 匹配一个或多个词,* 匹配不多不少恰好1个词,例如:item.# 能够匹配 item.insert.abc 或者 item.insert,item.* 只能匹配 item.insert
下面的只会发送给2
package com.ghx.topic;import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;import java.io.IOException;
import java.util.concurrent.TimeoutException;/*** @author :guo* @date :Created in 2025/3/20 11:35* @description:* @version:*/
public class Test01 {public static void main(String[] args) throws IOException, TimeoutException {//创建连接工厂ConnectionFactory factory=new ConnectionFactory();//rabbitmq服务器地址 默认本地localhostfactory.setHost("121.196.229.251");//端口号 默认5672factory.setPort(5672);//用户名 密码 默认guestfactory.setUsername("guest");factory.setPassword("guest");factory.setVirtualHost("/");//创建连接对象Connection connection=factory.newConnection();//获取channel对象Channel channel = connection.createChannel();//创建交换机
// String exchange,交换机的名称
// BuiltinExchangeType type, 交换机的类型
// boolean durable: 是否持久化channel.exchangeDeclare("topic_exchange", BuiltinExchangeType.TOPIC,true);//创建队列channel.queueDeclare("topic_queue1",true,false,false,null);channel.queueDeclare("topic_queue2",true,false,false,null);//绑定队列和交换机
// String queue,队列名
// String exchange,交换机名
// String routingKey: 路由key 因为广播模式没有路由key ""channel.queueBind("topic_queue1","topic_exchange","*.orange.*");channel.queueBind("topic_queue2","topic_exchange","*.*.rabbit");channel.queueBind("topic_queue2","topic_exchange","lazy.#");//发送消息String msg="hello topic交换机";channel.basicPublish("topic_exchange","lazy.orange",null,msg.getBytes());channel.close();connection.close();}
}