首页 分类

基本概念

AMQP协议:https://www.rabbitmq.com/tutorials/amqp-concepts.html
高级消息队列协议(Advanced Message Queue Protocol)
● 生产者:发消息到某个交换机
● 消费者:从某个队列中取消息
● 交换机(Exchange):负责把消息转发到对应的队列
● 队列(Queue):存储消息的
● 路由(Routes):转发,就是怎么把消息从一个地方转到另一个地方(比如从生产者转发到某个队列)

需要下载erlang作为依赖
https://www.erlang.org/patches/otp-27.0

启动

rabbitmq启动后,打开http://localhost:15672就可以打开rabbitmq的UI界面。如果不行,尝试着重启一下rabbitmq服务(用户名和密码都默认是guest)

注意

如果你打算在自己的服务器上部署RabbitMQ,你必须知道在默认设置下,使用gust账号是无法访问的。为何会这样呢?原因在于系统安全性的需求。如果你想从远程服务器访问管理面板,你需要创建一个新的管理员账号,不能使用默认的guest账号,否测会被系统拦截,导致无法访问或登泉。
点击右侧的User guest,会发现该用户拥有管理员权限,但这仅限于本地访问。如果你直接上线并使用默认的管理员权限,你可以想象可能出现的问题。比如,攻击者可能会扫描你服务器的端口,并不断尝试用guest账号登录。一旦找到一个开放的端口并用guest登录,攻击者就可以入侵你的服务器,然后在消息队列中不断添加消息,填满你的服务器硬盘。
因此,为了安全,系统默认关闭了gust账号的远程访问权限,这也是官方出于安全考虑的措施。关于如何创建管理员账号,在这里就不演示,可以参考官方文档进行操作一一官方文档的Adding a User。

快速入门
现在我们的RabbitMQ已经算启动了,去官网点击RabbitMQ Tutorials(RabbitMQ教程)。
引入消息队列java客户端
<!-- https://mvnrepository.com/artifact/com.rabbitmq/amqp-client -->
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.24.0</version>
</dependency>

java代码:
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

import java.nio.charset.StandardCharsets;

public class Send {

private final static String QUEUE_NAME = "hello";

public static void main(String[] argv) throws Exception {
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    try (Connection connection = factory.newConnection();
         Channel channel = connection.createChannel()) {
        channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        String message = "Hello World!";
        channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8));
        System.out.println(" [x] Sent '" + message + "'");
    }
}

}
channel(频道):理解为操作消息队列的client(比如jdbcClient、redisClient),提供了和消息队列server建立通信的传输方法。

创建消息队列:
参数:
queueName:消息队列名称(注意,同名称的消息队列,只能用同样的参数创建一次)
durabale:消息队列重启后,消息是否丢失
exclusive:是否只允许当前这个创建消息队列的连接操作消息队列
autoDelete:没有人用队列后,是否要删除队列

多消费者

https://www.rabbitmq.com/tutorials/tutorial-two-java.html

package com.yupi.springbootinit.mq;

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

public class MultiConsumer {
// 声明队列名称为"multi_queue"
private static final String TASK_QUEUE_NAME = "multi_queue";

public static void main(String[] argv) throws Exception {

// 创建一个新的连接工厂
ConnectionFactory factory = new ConnectionFactory();
// 设置连接工厂的主机地址
factory.setHost("localhost");
// 从工厂获取一个新的连接
final Connection connection = factory.newConnection();
// 从连接获取一个新的通道
final Channel channel = connection.createChannel();
// 声明一个队列,并设置属性:队列名称,持久化,非排他,非自动删除,其他参数;如果队列不存在,则创建它
channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null);
// 在控制台打印等待消息的提示信息
System.out.println(" [*] Waiting for messages. To exit press CTRL+C");
// (这个先注释)设置预取计数为1,这样RabbitMQ就会在给消费者新消息之前等待先前的消息被确认
// channel.basicQos(1);
  
// 创建消息接收回调函数,以便接收消息
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
    // 将接收到的消息转为字符串
    String message = new String(delivery.getBody(), "UTF-8");

    try {
        // (放到try里)打印出接收到的消息
            System.out.println(" [x] Received '" + message + "'");
        // 处理工作,模拟处理消息所花费的时间,机器处理能力有限(接收一条消息,20秒后再接收下一条消息)
        Thread.sleep(20000);
        // (不用doWork)模拟消息处理工作
        // doWork(message);
    } catch (InterruptedException e) {
        // 模拟处理消息所花费的时间
        e.printStackTrace();
    } finally {
        // 打印出完成消息处理的提示
        System.out.println(" [x] Done");
        // 手动发送应答,告诉RabbitMQ消息已经被处理
        channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
    }
};
// 开始消费消息,传入队列名称,是否自动确认,投递回调和消费者取消回调
channel.basicConsume(TASK_QUEUE_NAME, false, deliverCallback, consumerTag -> { });

}
// (不用doWork)用于模拟消息处理的函数,消息中的每一个'.'字符都会让线程暂停一秒钟
// private static void doWork(String task) {
// for (char ch : task.toCharArray()) {
// if (ch == '.') {
// try {
// // 暂停一秒钟
// Thread.sleep(1000);
// } catch (InterruptedException _ignored) {
// // 在被中断的情况下重新设置线程的中断状态
// Thread.currentThread().interrupt();
// }
// }
// }
// }
}
package com.yupi.springbootinit.mq;

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

public class MultiConsumer {
// 声明队列名称为"multi_queue"
private static final String TASK_QUEUE_NAME = "multi_queue";

public static void main(String[] argv) throws Exception {

// 创建一个新的连接工厂
ConnectionFactory factory = new ConnectionFactory();
// 设置连接工厂的主机地址
factory.setHost("localhost");
// 从工厂获取一个新的连接
final Connection connection = factory.newConnection();
// 从连接获取一个新的通道
final Channel channel = connection.createChannel();
// 声明一个队列,并设置属性:队列名称,持久化,非排他,非自动删除,其他参数;如果队列不存在,则创建它
channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null);
// 在控制台打印等待消息的提示信息
System.out.println(" [*] Waiting for messages. To exit press CTRL+C");
// (这个先注释)设置预取计数为1,这样RabbitMQ就会在给消费者新消息之前等待先前的消息被确认
// channel.basicQos(1);
  
// 创建消息接收回调函数,以便接收消息
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
    // 将接收到的消息转为字符串
    String message = new String(delivery.getBody(), "UTF-8");

    try {
        // (放到try里)打印出接收到的消息
            System.out.println(" [x] Received '" + message + "'");
        // 处理工作,模拟处理消息所花费的时间,机器处理能力有限(接收一条消息,20秒后再接收下一条消息)
        Thread.sleep(20000);
    } catch (InterruptedException e) {
        // 模拟处理消息所花费的时间
        e.printStackTrace();
    } finally {
        // 打印出完成消息处理的提示
        System.out.println(" [x] Done");
        // 手动发送应答,告诉RabbitMQ消息已经被处理
        channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
    }
};
// 开始消费消息,传入队列名称,是否自动确认,投递回调和消费者取消回调
channel.basicConsume(TASK_QUEUE_NAME, false, deliverCallback, consumerTag -> { });

}
}

消息确认机制

这里有个很有意思的地方,叫autoack,默认为false一消息确认机制。
消息队列如何确保消费者已经成功取出消息呢?它依赖一个称为消息确认的机制。当消费者从队列中取走消息后,必须对此进行确认。这就像在收到快递后确认收货一样,这样消息队列才能知道消费者已经成功取走了消息,并能安心地停止传输。因此,整个过程就像这样。
在面试中,如果被问及如何保证消息不会丢失,例如当业务流程失败时该怎么办,你可以回答:“我们可以选择拒绝接收失败的消息,并重新启动。重新启动或进行其他处理可以指定拒绝某条消息。”

为了保证消息成功被消费(快递成功被取走),rabbitmg提供了消息确认机制,当消费者接收到消息后,比如要给一个反馈:
● ack:消费成功
● nack:消费失败
● reject:拒绝
如果告诉rabbitmg服务器消费成功,服务器才会放心地移除消息。
channel.basicConsume(TASK_QUEUE_NAME, false, deliverCallback, consumerTag -> {});
但是,如果在接收到消息后,工作尚未完成,我们是否就不需要确认成功呢?这种情况,建议将autoack设置为false,根据实际情况手动进行确认了。
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
第二个参数'multiple'表示比量确认,也就是说,是否需要一次性确认所有的历史消息,直到当前这条消息为止。
第3个参数表示是否重新入队,可用于重试。

交换机

一个生产者给多个队列发消息,1个生产者对多个队列。
交换机的作用:提供消息转发功能,类似于网络路由器
要解决的问题:怎么把消息转发到不同的队列上,好让消费者从不同的队列消费。
教程:https://www.rabbitmq.com/tutorials/tutorial-three-java.html
交换机有多种类别:fanout、direct、topic、headers and.

fanout

扇出、广播
特点:消息会被转发到所有绑定到该交换机的队列
场景:很适用于发布订阅的场景。比如写日志,可以多个系统间共享
注意:
1.消费者和生产者要绑定同一个交换机
2.要先有队列,才能绑定
一对示例:

package com.cy.cyBIbackend.mq;

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

import java.util.Scanner;

public class FanoutProducer {

private static final String EXCHANGE_NAME = "fanout-exchange";

public static void main(String[] argv) throws Exception {

ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
     Channel channel = connection.createChannel()) {

    // 创建交换机
    channel.exchangeDeclare(EXCHANGE_NAME, "fanout");

    Scanner scanner = new Scanner(System.in);
    while(scanner.hasNext()){
        String message = scanner.nextLine();
        channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes("UTF-8"));
        System.out.println(" [x] Sent '" + message + "'");

    }
}

}

}

package com.cy.cyBIbackend.mq;

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

public class FanoutConsumer {
private static final String EXCHANGE_NAME = "fanout-exchange";

public static void main(String[] argv) throws Exception {

ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection();
Channel channel1 = connection.createChannel();
Channel channel2 = connection.createChannel();
// 声明交换机
channel1.exchangeDeclare(EXCHANGE_NAME, "fanout");
// 创建队列,随机分配一个队列名称
String queueName = "小王的工作队列";
channel1.queueDeclare(queueName,true,false, false,null);
channel1.queueBind(queueName, EXCHANGE_NAME, "");

String queueName2 = "小李的工作队列";
channel2.queueDeclare(queueName2,true,false, false,null);
channel2.queueBind(queueName2, EXCHANGE_NAME, "");

System.out.println(" [*] Waiting for messages. To exit press CTRL+C");

DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
    String message = new String(delivery.getBody(), "UTF-8");
    System.out.println(" [小王] Received '" + message + "'");
};

DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
  String message = new String(delivery.getBody(), "UTF-8");
  System.out.println(" [小李] Received '" + message + "'");
};

channel1.basicConsume(queueName, true, deliverCallback1, consumerTag -> { });
channel2.basicConsume(queueName2, true, deliverCallback2, consumerTag -> { });

}
}

Direct交换机

绑定:可以让交换机和队列进行关联,可以指定让交互机把什么样的消息发送给哪个队列
routingKey:路由键,控制消息要转发给哪个队列的
特点:消息会根据路由键转发到指定的队列
场景:特定的消息只交给特定的系统(程序)来处理
一个示例:

package com.cy.cyBIbackend.mq;

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

import java.util.Scanner;

public class DirectProducer {

private static final String EXCHANGE_NAME = "direct_exchange";

public static void main(String[] argv) throws Exception {

ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
     Channel channel = connection.createChannel()) {
    channel.exchangeDeclare(EXCHANGE_NAME, "direct");

    Scanner scanner = new Scanner(System.in);
    while(scanner.hasNext()){
        String userInput = scanner.nextLine();
        String[] strings = userInput.split(" ");
        if(strings.length < 1){
            continue;
        }
        String message = strings[0];
        String routingKey = strings[1];
        channel.basicPublish(EXCHANGE_NAME, routingKey, null, message.getBytes("UTF-8"));
        System.out.println(" [x] Sent '" + message + "'with routing:'" + message + "'");

    }
}

}
//..
}

package com.cy.cyBIbackend.mq;

import com.rabbitmq.client.*;

public class DirectConsumer {

private static final String EXCHANGE_NAME = "direct_exchange";

public static void main(String[] argv) throws Exception {

ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();

channel.exchangeDeclare(EXCHANGE_NAME, "direct");

String queueName1 = "xiaochen_queue";
channel.queueDeclare(queueName1,true,false, false,null);
channel.queueBind(queueName1, EXCHANGE_NAME, "xiaochen");
String queueName2 = "xiaoyang_queue";
channel.queueDeclare(queueName2,true,false, false,null);
channel.queueBind(queueName2, EXCHANGE_NAME, "xiaoyang");

// String queueName = channel.queueDeclare().getQueue();

System.out.println(" [*] Waiting for messages. To exit press CTRL+C");

DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
    String message = new String(delivery.getBody(), "UTF-8");
    System.out.println(" [xiaochen] Received '" +
        delivery.getEnvelope().getRoutingKey() + "':'" + message + "'");
};
  DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
      String message = new String(delivery.getBody(), "UTF-8");
      System.out.println(" [xiaoyang] Received '" +
              delivery.getEnvelope().getRoutingKey() + "':'" + message + "'");
  };
channel.basicConsume(queueName1, true, deliverCallback1, consumerTag -> { });
channel.basicConsume(queueName2, true, deliverCallback2, consumerTag -> { });

}
}

topic交换机

特点:消息会根据一个模糊的路由键转发到指定的队列
场景:特定的一类消息可以交给特定的一类系统(程序)来处理
绑定关系:可以模糊匹配多个绑定
:匹配一个单词,比如.orange,那么a.orange、b.orange都能匹配
● #:匹配0个或多个单词,比如a.#,那么a.a、a.b、a.a.a都能匹配

RPC

支持用消息队列来模拟RPC的调用,但是一般没必要,直接用Dubbo、GRPC等RPC框架就好了。

适用场景:清理过期数据、模拟延迟队列的实现(不开会员就慢速)、专门让某个程序处理过期请求

● 消息过期机制是干嘛的?消息过期机制是用来处理那些在一段时间内未被处理的消息。颀名思义,当一条消息在一定时间内未被消费者处理时,它就会过期,即失效。这种机制允许系统自动清理和丢弃那些长时间未被消费的消息,以避免消息队列中积累过多的过期消息,从而保持系统的效率和可靠性。

● 消息过期机制有什么应用场景?举个例子,假设一个用户向一个新系统发送了一个请求,比如支付订单。订单通常有一个有效期限,比如15分钟。如果在这15分钟内,下游系统没有及时处理该订单,这可能意味着下游系统出现了故障,或者在这个时间段内不再需要处理该订单了。在这种情况下,如果将该消息放入队列中,并设置了15分钟的过期时间,那么如果15分钟后还没有消费者来获取该消息,该支付订单消息就会过期。订单在15分钟内没有支付,那么该订单就已经失效了。下游系统就不再需要处理这个订单,也不需要更新库存或处理物流等相关操作。因此,消息过期机制非常适用于这种过期场景的处理。通过设置合适的过期时间,可以确保及时清理无效的消息,提高系统的效率和准确性。

● 什么叫延迟队列呢?它允许将消息延迟一段时间后再进行处理。举个例子,假设我们有一个程序,要求在用户进行某项操作之后不是立即处理,而是延迟几分钟后再执行。这种情况下,延迟队列可以派上用场。让我们以区分普通用户和会员用户的场景为例。对于会员用户,我们希望立即处理其请求;而对于普通用户,我们希望让其排队等待一段时间(比如5分种)后再进行处理,以鼓励其购买会员。这时,可以利用延迟队列实现。消费者可以监听延迟队列,普通用户的请求由一个程序处理监听该延迟队列,而会员用户的请求则由另一个程序监听一个高优先级的队列。一旦你掌握了消息队列的知识,就可以实现这样的程序逻辑。延迟队列的实现可以借助消息过期机制。具体的实现思路是创建两个队列,第一个队列中的消息设置了过期时间,比如5分钟,然后将过期的消息转移到第二个队列中。接着,让相应的用户程序监听第二个队列,这样第二个队列就成为了延迟队列。实际上,这个思路与我们将要讨论的死信队列有一些相似之处,因为它们都涉及到特定的场景和处理方式。

rabbitMQ相当于是两种消息过期机制:
● 第一种方式是给队列中的所有消息指定一个统一的过期时间。也就是说,无论何时进入这个队列的消息,在特定的时间点都会过期失效。这种方式是针对整个队列而言。
● 第二种方式是给某条具体的消息指定过期时间。这意味着,针对特定的消息,我们可以指定一个独立的过期时间。这样,在达到指定的时间后,这条消息将会过期并自动失效。
● 这两种方式在应用中具有不同的应用场景:第一种方式适用于需要在一定时间后对整个队列中的消息进行处理或清理的场景。例如,我们可以设置一个定期的清理任务,删除队列中过期的消息,以确保队列的有效性和性能。而第二种方式则适用于对于某些特定消息需要具备独立过期时间的场景。比如,在电商平台中,如果有一个购物车消息,我们可以为每个购物车消息单独设置过期时间,以便在一段时间内保留购物车状态,如果超过过期时间,系统可以自动清理过期的购物车消息。

设置代码:

Map<String, Object> args = new HashMap<String, Object>();

args.put("x-message-ttl", 60000);
channel.queueDeclare("myqueue", false, false, false, args);

package com.yupi.springbootinit.mq;

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

import java.nio.charset.StandardCharsets;

public class TtlProducer {

// 定义队列名称为"ttl_queue"
private final static String QUEUE_NAME = "ttl_queue";

public static void main(String[] argv) throws Exception {
    // 创建连接工厂
    ConnectionFactory factory = new ConnectionFactory();
    factory.setHost("localhost");
    // 建立连接、创建频道
    try (Connection connection = factory.newConnection();
         Channel channel = connection.createChannel()) {
        // 消息虽然可以重复声明,必须指定相同的参数,在消费者的创建队列要指定过期时间,
        // 后面要放args,在生产者你又想重新创建队列,又不指定参数,那肯定会有问题,
        // 所以要把这里的创建队列注释掉。
        // channel.queueDeclare(QUEUE_NAME, false, false, false, null);
        
        // 发送消息
        String message = "Hello World!";
        // 使用默认的交换机,将消息发送到指定队列
        channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8));
        System.out.println(" [x] Sent '" + message + "'");
    }
}

}

package com.yupi.springbootinit.mq;

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

import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;

public class TtlConsumer {

// 定义我们正在监听的队列名称"ttl_queue"
private final static String QUEUE_NAME = "ttl_queue";

public static void main(String[] argv) throws Exception {
    // 创建连接工厂
    ConnectionFactory factory = new ConnectionFactory();
    // 设置连接工厂的主机地址为 "localhost"
    factory.setHost("localhost");
    // 建立连接
    Connection connection = factory.newConnection();
    // 创建频道
    Channel channel = connection.createChannel();

    // 创建队列,指定消息过期参数
    Map<String, Object> args = new HashMap<String, Object>();
    // 设置消息过期时间为5秒
    args.put("x-message-ttl", 5000);
    // 创建队列,并传入队列名称、是否持久化、是否私有、是否自动删除,args 指定参数
    channel.queueDeclare(QUEUE_NAME, false, false, false, args);
    // 打印等待消息的提示信息
    System.out.println(" [*] Waiting for messages. To exit press CTRL+C");
    // 定义了如何处理消息的回调函数
    DeliverCallback deliverCallback = (consumerTag, delivery) -> {
        String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
        System.out.println(" [x] Received '" + message + "'");
    };
    // 消费消息,该方法会持续阻塞,等待接收消息
    channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> { });
}

}

如果在过期时间内,还没有消费者取消息,消息才会过期。
注意,如果消息已经接收到,但是没确认,是不会过期的。

指定过期时间:
// 给消息指定过期时间
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()

    // 设置消息的过期时间为1000毫秒
    .expiration("1000")
    .build();

// 发布消息到指定的交换机("my-exchange")和路由键("routing-key")
// 使用指定的属性(过期时间)和消息内容(UTF-8编码的字节数组)
channel.basicPublish("my-exchange", "routing-key", properties, message.getBytes(StandardCharsets.UTF_8));

死信队列

为了保证消息的可靠性,比如每条消息都成功消费,需要提供一个容错机制,即:失败的消息怎么处理?
死信:过期的消息、拒收的消息、处理失败的消息的统称
死信队列:专门处理死信的队列(注意,它就是一个普通队列,只不过是专门用来处理死信的,你甚至可以理解这个队列的名称叫“死信队列”)
死信交换机:专门给死信队列转发消息的交换机(注意,它就是一个普通交换机,只不过是专门给死信队列发消息而已,理解为这个交换机的名称就叫“死信交换机”)
死信可以通过死信交换机绑定到死信队列。
实现:
1.创建死信交换机和死信队列,并且绑定关系
2.给失败之后需要容错处理的队列绑定死信交换机



文章评论