rabbitmq 从字面意思来看就是兔子消息队列,总之就是像兔子一样快并且繁殖能力强

mq 的作用

1.异步解耦 2.流量削峰 3.消息分发 4.延迟通知

市面上最主流的 mq 产品

1.Kafka

Kafka 适用于对 日志 收集和传输     有需求的场景,它追求高吞吐量,性能卓越,主要支持简单的 MQ 功能,日志领域比较成熟

2.RocketMQ

使用 java 语言开发,在可用性,可靠性,稳定性等方面比较出色,适合对于可靠性比较高,且并发比较大的场景

3.RabbitMQ

采用 Erlang 语言开发,MQ 功能比较完备,几乎支持所有的主流语言,性能较好,吞吐量能达到万级,适合数据量没那么大,并发量没那么高的场景

  

我们理解这六部分内容就清楚 rabbitmq 到底是怎样工作的了

RabbitMQ 是一个消息中间件,也是一个生产者消费者模型,他负责接收存储并转发消息

Producer 和 Comsumer(生产者和消费者)

这么说吧生产者负责生产消息,它还会额外携带着一些特殊的字符串,用来把自己推销给那么喜欢自己的消费者,消费者负责消费消息,它只会收到消息,而并不知道是谁发送的。

Broker:

RabbitMQ Server,负责接收和发送消息,可以把他简单的看作是一个 RabbitMQ 服务实例

Connection 和 Channel

Connection:

连接,是客户端和RabbitMQ服务器之间的一个TCP连接,这个连接时建立消息传递的基础,它负责传输客户端和服务器之间的所有数据和控制信息

Channel:

通道,在 RabbitMQ 中,一个 TCP 连接可以有多个 Channel,每个 Channel 都是独立的虚拟连接,消息的发送和接收都是基于 Channel

通道的主要作用是将消息的读写操作复用到同一个 TCP 连接上,这样可以减少建立和关闭连接的开销,提高性能

Virtual host (虚拟主机)

它为消息队列提供了一种逻辑上的隔离机制,对于 RabbitMQ 来说,一个 BrokerServer 上可以存在多个 Virtual Host,当多个不同的用户使用同一个 RabbitMQ Server 提供服务时,可以划分出多个 vhost,每个用户在自己的 vhost 创建 交换机 和 队列

Queue (队列)

作为 RabbitMQ 的内部对象,用于存储消息

Exchange(交换机)

它负责接收生产者发送的消息,并根据特定的规则把这些消息路由到队列中

现在重新来看一下这个

工作流程图

1.Producer 生产了一条消息
2.Producer 连接到 RabbitMQBroker,建立一个连接,开启一个通道
3.Producer 声明一个交换机,路由消息
4.Producer 声明一个队列,存放消息
5.Producer 发送消息到 RabbitMQ Broker
6.RabbitMQ Broker 接收消息,并存入相应的队列中,如果未找到相应的队列,则会根据生产者配置,丢弃或者退回

AMQP

AMQP 是一种高级消息队列协议,它定义了一套确定的消息交换功能,使得生产者能够将消息发送到交换器,然后由队列接收等待消费者消费

我们来测试一下这个功能

引入依赖
        <dependency>
            <groupId>com.rabbitmq</groupId>
            <artifactId>amqp-client</artifactId>
            <version>5.20.0</version>
        </dependency>
编写生产者
1.创建连接
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost(Constants.HOST);
        connectionFactory.setPort(Constants.PORT);// 需要提前开放端口号
        connectionFactory.setUsername(Constants.USER_NAME);// 账号
        connectionFactory.setPassword(Constants.PASSWORD);// 密码
        connectionFactory.setVirtualHost(Constants.VIRTUAL_HOST);// 虚拟机
        Connection connection = connectionFactory.newConnection();
2.创建通道
        Channel channel = connection.createChannel();
3.声明一个队列
        channel.queueDeclare("hello",true,false,false,null);

这五个参数的含义

1.队列的名称     2.是否持久化    3.是否只让一个消费者监听队列,在连接关闭时,是否删除队列      4.是否自动删除    5.传入参数

4.发送消息

当一个新的 RabbitMQ 节点启动时,他会预先声明几个内置的交换机,内置交换机是空字符串,生产者发送的消息会根据队列名称直接路由到对应的队列

 for (int i = 0; i < 10; i++) {
            String msg = "hello rabbitmq:" + i;
            channel.basicPublish("","hello",null,msg.getBytes());
        }

这几个参数的含义

1.交换机名称

2.路由名称

3.配置信息

4.发送消息的数据

5.释放资源

// 资源释放
        channel.close();
        connection.close();
执行代码

我们发现了声明的队列

消息也确实被放在了队列中了

package rabbitmq.simple;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

public class ProducerDemo {
    public static void main(String[] args) throws IOException, TimeoutException {
        // 建立连接
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("");
        connectionFactory.setPort(5672);// 需要提前开放端口号
        connectionFactory.setUsername("study");// 账号
        connectionFactory.setPassword("study");// 密码
        connectionFactory.setVirtualHost("");// 虚拟机
        Connection connection = connectionFactory.newConnection();

        // 开启信道
        Channel channel = connection.createChannel();
        // 声明交换机     使用内置的交换机
        // 声明队列
        /**
         *  Queue.DeclareOk queueDeclare(String queue, boolean durable, boolean exclusive, boolean autoDelete,
         *                                  Map<String, Object> arguments) throws IOException;
         *  参数说明
         *  queue: 队列名称
         *  durable: 可持久化
         *  exclusive: 是否独占
         *  autoDelete: 是否自动删除
         *  arguments:
         */
        channel.queueDeclare("hello",true,false,false,null);
        // 发送消息
        /**
         * basicPublish(String exchange, String routingKey, BasicProperties props, byte[] body)
         * 参数说明:
         * exchange: 交换机名称
         * rotingKey: 内置交换机,routingKey 和 队列名称保持一致
         * props: 属性配置
         * body: 消息
         */
        for (int i = 0; i < 10; i++) {
            String msg = "hello rabbitmq:" + i;
            channel.basicPublish("","hello",null,msg.getBytes());
        }
        // 资源释放
        channel.close();
        connection.close();
    }
}
编写消费者

1.创建连接 2.创建通道 3.声明队列 4.消费消息 5.释放资源

消费当前队列
DefaultConsumer consumer = new DefaultConsumer(channel){
            // 从队列中接收到消息,执行的方法
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                System.out.println("接收到消息:" + new String(body));
            }
        };
        channel.basicConsume("hello",true,consumer);
释放资源
 // 等待程序执行完成
        Thread.sleep(2000);
        // 释放资源
        channel.close();
        connection.close();

package rabbitmq.simple;

import com.rabbitmq.client.*;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

public class ConsumerDemo {
    public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
        // 创建连接
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setHost("");
        connectionFactory.setPort(5672);// 需要提前开放端口号
        connectionFactory.setUsername("study");// 账号
        connectionFactory.setPassword("study");// 密码
        connectionFactory.setVirtualHost("");// 虚拟机
        Connection connection = connectionFactory.newConnection();
        // 创建 channel
        Channel channel = connection.createChannel();
        // 声明队列,
        channel.queueDeclare("hello",true,false,false,null);
        // 消费消息
        /**
         *  String basicConsume(String queue, boolean autoAck, Consumer callback)
         *  参数说明
         *  queue: 队列名称
         *  autoAck: 是否自动确认
         *  callback: 接收到消息后,执行的逻辑
         */
        DefaultConsumer consumer = new DefaultConsumer(channel){
            // 从队列中接收到消息,执行的方法
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                System.out.println("接收到消息:" + new String(body));
            }
        };
        channel.basicConsume("hello",true,consumer);
        // 等待程序执行完成
        Thread.sleep(2000);
        // 释放资源
        channel.close();
        connection.close();

    }
}

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐