当前位置: 首页 > article >正文

MQ,RabbitMQ,MQ的好处,RabbitMQ的原理和核心组件,工作模式

1.MQ

 MQ全称 Message Queue(消息队列),是在消息的传输过程中  保存消息的容器。它是应用程序和应用程序之间的通信方法

1.1 为什么使用MQ

在项目中,可将一些无需即时返回且耗时的操作提取出来,进行异步处理,而这种异步处理的方式大大的节省了服务器的请求响应时间,从而提高系统吞吐量

1.2MQ的好处

1.应用解耦   系统间通过消息通信,不用关心其他系统的处理。

2.异步提速  相比于传统的串行、并行方式,提高了系统吞吐量。

3.削峰填谷   可以通过消息队列长度控制请求量;可以缓解短时间内的高并发请求。

简单来说: 就是在访问量剧增的情况下,但是应用仍然不能停,比如“双十一”下单的人多,但是淘宝这个应用仍然要运行,所以就可以使用消息中间件采用队列的形式减少突然访问的压力

使用MQ后,可以提高系统稳定性

1.3劣势

  1. 系统可用性降低 系统引入的外部依赖越多,系统稳定性越差。一旦 MQ 宕机,就会对业务造成影响。如何保证MQ的高可用?

  2. 系统复杂度提高 MQ 的加入大大增加了系统的复杂度,以前系统间是同步的远程调用,现在是通过 MQ 进行异步调用。如何保证消息没有被重复消费?怎么处理消息丢失情况?那么保证消息传递的顺序性?

  3. 一致性问题 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的原理

核心组件

  1. 生产者(Producer):负责发送消息到交换器的客户端应用程序。

  2. 消费者(Consumer):从队列中获取并处理消息的客户端应用程序。

  3. 交换器(Exchange):接收生产者发送的消息,并根据路由规则将消息转发到相应的队列。

  4. 队列(Queue):存储消息,直到消费者取走消息。

  5. 绑定(Binding):定义交换器和队列之间的关联关系。

工作流程

  1. 消息发送:生产者通过信道(Channel)将消息发送到交换器。

  2. 消息路由:交换器根据路由键(Routing Key)和绑定键(Binding Key)将消息路由到相应的队列。

  3. 消息存储:队列存储消息,等待消费者取走。

  4. 消息消费:消费者通过信道从队列中获取消息并处理。

交换器类型

  1. Direct:根据完全匹配的路由键将消息发送到相应的队列。

  2. Fanout:将消息广播到所有绑定的队列,不考虑路由键。

  3. 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服务器地址 默认本地localhost
        factory.setHost("xxxx");
        //端口号 默认5672
        factory.setPort(5672);
        //用户名 密码  默认guest
        factory.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服务器地址 默认本地localhost
        factory.setHost("xxxx");
        //端口号 默认5672
        factory.setPort(5672);
        //用户名 密码  默认guest
        factory.setUsername("guest");
        factory.setPassword("guest");
        //虚拟机名称 默认/
        factory.setVirtualHost("/");
        //创建连接对象
        Connection connection = factory.newConnection();
        //获取channel对象
        Channel channel = connection.createChannel();
        DefaultConsumer consumer=new DefaultConsumer(channel){
            @Override
            public 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){
            @Override
            public 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){
            @Override
            public 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种类型:

  1. Fanout:广播,将消息交给所有绑定到交换机的队列

  2. Direct:定向,把消息交给符合指定routing key 的队列

  3. 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服务器地址 默认本地localhost
        factory.setHost("xxxx");
        //端口号 默认5672
        factory.setPort(5672);
        //用户名 密码  默认guest
        factory.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服务器地址 默认本地localhost
        factory.setHost("xxxx");
        //端口号 默认5672
        factory.setPort(5672);
        //用户名 密码  默认guest
        factory.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服务器地址 默认本地localhost
        factory.setHost("121.196.229.251");
        //端口号 默认5672
        factory.setPort(5672);
        //用户名 密码  默认guest
        factory.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();


    }
}


http://www.kler.cn/a/593689.html

相关文章:

  • 飞行控制系统SRAM存储解决方案
  • 测试详解 (概念篇、Bug篇、用例篇、测试分类)
  • wpa_supplicant二层包收发
  • Gone gRPC 组件使用指南
  • Android BLE 权限管理
  • 论文阅读 EEGNet
  • Vue:添加响应式数据
  • 使用git托管项目
  • 【DeepSeekR1】怎样清除mssql的日志文件?
  • 服务的拆分数据的迁移
  • ambiq apollo3 Flash实例程序注释
  • Java 实现两个线程交替打印AB的几种方式
  • 从零开始搭建向量数据库:基于 Xinference 和 Milvus 的文本搜索实践
  • 机器学习是怎么一步一步由神经网络发展到今天的Transformer架构的?
  • 2025 使用docker部署ubuntu24容器并且需要ubuntu24容器能通过ssh登录SSH 登录的Ubuntu24容器
  • Modern C++处理 Hooks 机制
  • Datawhale大语言模型-Transformer以及模型详细配置
  • HttpClient通讯时间过久
  • MiniMax GenAI 可观测性分析:基于阿里云 SelectDB 构建 PB 级别日志系统
  • python采集小红书笔记详情API接口,json数据示例分享