了解 RabbitMQ
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();
}
}
更多推荐




所有评论(0)