欢迎光临
我们一直在努力

RabbitMQ是什么?如何使用

1. 什么是MQ

消息队列(Message Queue,简称MQ)

  • 从字面意思上看,本质是个队列,FIFO先入先出,只不过队列中存放的内容是message而已。 其主要用途:不同进程Process/线程Thread之间通信。

为什么会产生消息队列?有几个原因:

  • 不同进程(process)之间传递消息时,两个进程之间耦合程度过高,改动一个进程,引发必须修改另一个进程,为了隔离这两个进程,在两进程间抽离出一层(一个模块),所有两进程之间传递的消息,都必须通过消息队列来传递,单独修改某一个进程,不会影响另一个;

  • 不同进程(process)之间传递消息时,为了实现标准化,将消息的格式规范化了,并且,某一个进程接受的消息太多,一下子无法处理完,并且也有先后顺序,必须对收到的消息进行排队,因此诞生了事实上的消息队列;

MQ框架非常的多,比较流行的有RabbitMq、kafka,以及阿里开源的RocketMQ。本文主要介绍RabbitMq。

区别

优点缺点适用场景
kafka 吞吐量非常大,性能非常好,技术生态完整 功能比较单一 分布式日志搜集、大数据采集,linkin用来处理分布式日志用的
RabbitMQ 消息可靠性强,功能全面 吞吐量较低。消息挤压会影响性能,erlang语言比较小众 企业内部系统调用
RocketMQ 高吞吐、高性能、高可用、高级功能非常全 技术生态不是很完整 几乎全场景,尤其适合金融 ,阿里拿来做金融用的

2. 为什么使用消息队列

主要有三个作用:

  • 解耦。如图所示。本身生产者所需要的数据被消费者们所需要,这要给他们都添加响应的方法,但如果消费者3突然不需要该信息,生产者就要删掉该方法,但如果增加这个中间层进行解耦,就可以让他们只去找这个消息队列去拿内容,而不需要再次去寻找生产者,生产者可以更专注于自己的业务。 在这里插入图片描述

  • 异步。如图所示。一个客户端请求发送进来,系统A会调用系统B、C、D三个系统,同步请求的话,响应时间就是系统A、B、C、D的总和,也就是800ms。如果使用MQ,系统A发送数据到MQ,然后就可以返回响应给客户端,不需要再等待系统B、C、D的响应,可以大大地提高性能。对于一些非必要的业务,比如发送短信,发送邮件等等,就可以采用MQ。 在这里插入图片描述

  • 削峰。就是当业务量大的时候,让他做一个桥的作用,不让数据一股脑打到数据库上,让数据库宕机。使用MQ,是让sql语句不在直接打到数据库,而是把数据发送到MQ,MQ短时间积压数据是可以接受的,然后由消费者每次拉取2000条进行处理,防止在请求峰值时期大量的请求直接发送到MySQL导致系统崩溃。

3. RabbitMQ介绍

RabbitMQ 是一个开源的消息代理和队列服务器,用于在分布式系统中存储、转发和接收消息。它是基于高级消息队列协议(AMQP)实现的,支持多种客户端和协议,可以用于多种场景,如负载均衡、分布式事务处理、消息通知等。

RabbitMQ 的核心概念

RabbitMQ 的核心概念包括生产者、消费者、队列、交换机和绑定。生产者负责发送消息到交换机,交换机根据路由规则将消息转发到绑定的队列,消费者从队列中获取消息进行处理。RabbitMQ 支持多种类型的交换机,如直接交换机(direct)、扇形交换机(fanout)、主题交换机(topic)和头交换机(headers),它们各自适用于不同的路由策略和模式。

RabbitMQ 的工作原理

RabbitMQ 的工作原理是通过Broker(消息代理服务器)来接收、存储和转发消息。Broker 包含一个或多个虚拟主机,每个虚拟主机可以有自己的队列、交换机和绑定。消息的生产者将消息发送到交换机,交换机根据绑定规则将消息路由到队列,消费者监听队列并处理消息。

RabbitMQ 的使用场景

RabbitMQ 可以用于实现系统间的解耦、异步处理和流量削峰。例如,在电商系统中,订单生成后可以将订单信息发送到消息队列,库存系统和物流系统可以从队列中获取订单信息进行处理,这样即使某个系统暂时不可用,也不会影响整个流程的进行。此外,RabbitMQ 还可以用于实现延迟消息和定时任务,如订单超时未支付自动取消等功能。

4. 使用

RabbitMQ是由ErLang开发的,他需要erlang的环境来进行运行。 所以要先安装erlang,你可以根据rabbitMQ官网来看erlang适应的版本。 在这里插入图片描述 安装好erlang之后,将erlang配置成环境变量,也就是将erlang的sbin目录保存在环境变量的path中。

然后官网安装windows包 在这里插入图片描述 这是安装好的目录。 在这里插入图片描述

安装完之后,他会自动执行,这里我出现了一个错误,就是他自动运行之后,并没有运行起来,我怀疑是windows的运行方式和其他的不一样。 那我是如何解决的呢,就是先通过windows的services.msc进入服务,关闭自动运行的rabbitMQ,然后调整成手动运行,之后在执行rabbitMQ的bat命令 运行命令

rabbitmqserver start

当你看见下面情况时候 在这里插入图片描述 证明运行成功了。 然后你可以访问web页面,http://localhost:15672 账号密码默认是:guest/guest

进入到下面页面,就大功告成了。 在这里插入图片描述

用户

查看当前拥有用户

rabbitmqctl list_users

查看权限

rabbitmqctl list_permissions

添加用户

rabbitmqctl add_user user_name password

设置用户tag

rabbitmqctl set_user_tags user_name administrator

设置用户权限

rabbitmqctl set_permissions p "/" admin ".*" ".*" ".*"

5. java操作

当然你的服务配置好,并没有结束,而是刚刚开始,我们需要用编程语言去操作他,这里我使用java。 首先引入环境,springboot是拥有mq环境的,所以只需要引入即可。

<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

依然是在配置yml中加入RabbitMQ的配置信息

spring:
rabbitmq:
host: 127.0.0.1
port: 5672
username: guest
password: guest

然后是创建RabbitMQ的配置类,放入到IOC当中。

@Configuration
public class DirectRabbitConfig {
@Bean
public Queue rabbitmqDemoDirectQueue() {
/**
* 1、name: 队列名称
* 2、durable: 是否持久化
* 3、exclusive: 是否独享、排外的。如果设置为true,定义为排他队列。则只有创建者可以使用此队列。也就是private私有的。
* 4、autoDelete: 是否自动删除。也就是临时队列。当最后一个消费者断开连接后,会自动删除。
* */

return new Queue(RabbitMQConfig.RABBITMQ_DEMO_TOPIC, true, false, false);
}

@Bean
public DirectExchange rabbitmqDemoDirectExchange() {
//Direct交换机
return new DirectExchange(RabbitMQConfig.RABBITMQ_DEMO_DIRECT_EXCHANGE, true, false);
}

@Bean
public Binding bindDirect() {
//链式写法,绑定交换机和队列,并设置匹配键
return BindingBuilder
//绑定队列
.bind(rabbitmqDemoDirectQueue())
//到交换机
.to(rabbitmqDemoDirectExchange())
//并设置匹配键
.with(RabbitMQConfig.RABBITMQ_DEMO_DIRECT_ROUTING);
}
}

然后是你运行操作数据所需要的服务类

@Service
public class RabbitMQServiceImpl implements RabbitMQService {
//日期格式化
private static SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");

@Resource
private RabbitTemplate rabbitTemplate;

@Override
public String sendMsg(String msg) throws Exception {
try {
String msgId = UUID.randomUUID().toString().replace("-", "").substring(0, 32);
String sendTime = sdf.format(new Date());
Map<String, Object> map = new HashMap<>();
map.put("msgId", msgId);
map.put("sendTime", sendTime);
map.put("msg", msg);
rabbitTemplate.convertAndSend(RabbitMQConfig.RABBITMQ_DEMO_DIRECT_EXCHANGE, RabbitMQConfig.RABBITMQ_DEMO_DIRECT_ROUTING, map);
return "ok";
} catch (Exception e) {
e.printStackTrace();
return "error";
}
}
}

削峰填谷

在高并发场景下,例如秒杀活动,用户请求量可能瞬间激增,导致服务器无法承受。RabbitMQ 通过消息队列的方式,将瞬时流量缓冲到队列中,然后以稳定的速率处理这些请求,从而实现削峰填谷。

配置文件示例

以下是一个 Spring Boot 项目的 RabbitMQ 配置文件示例:

spring.application.name=springboot_rabbitmq
spring.rabbitmq.host=192.168.0.102
spring.rabbitmq.port=5672
spring.rabbitmq.username=admin
spring.rabbitmq.password=admin
spring.rabbitmq.virtualhost=/
spring.rabbitmq.listener.simple.acknowledgemode=manual # 设置手动应答
spring.rabbitmq.listener.simple.prefetch=2 # 每次最多可处理信息量

消费者代码示例

消费者代码通过 @RabbitListener 注解监听队列,并手动应答消息:

@Component
@RabbitListener(queuesToDeclare = @Queue(name = "springboot-limit"))
public class CurrentlimitCustomer {
@RabbitHandler
public void receive(String msg, Channel channel, Message message) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
Thread.sleep(1000 * 10); // 模拟处理时间
System.out.println("=====限流====>");
System.out.println(msg);
System.out.println(channel);
System.out.println(message);
// 手动签收消息
channel.basicAck(deliveryTag, true);
} catch (Exception e) {
// 处理异常
}
}
}

生产者代码示例

生产者代码通过 rabbitTemplate 发送消息到队列:

@Test
public void test09() throws Exception {
for (int i = 0; i < 10; i++) {
rabbitTemplate.convertAndSend("springboot-limit", "限流测试");
}
Thread.sleep(1000 * 1000); // 模拟延迟
}

削峰填谷的优势

  • 提高系统稳定性:通过缓冲瞬时流量,避免系统过载。
  • 提升用户体验:减少请求失败的概率,提高系统响应速度。
  • 简化系统设计:通过消息队列实现异步处理,降低系统耦合度。
  • RabbitMQ 的削峰填谷功能在处理高并发场景中具有显著优势,能够有效提高系统的稳定性和可用性。

    交换机

    Exchange交换机是消息路由的核心组件,主要用于接收生产者发送的消息,并根据路由键(Routing Key)将消息分发到绑定的队列(Queue)。它在消息队列系统(如RabbitMQ)中起到类似“邮局”的作用,负责将消息分拣到正确的队列,从而实现生产者与消费者的解耦。

    交换机的核心功能

    生产者将消息发送到交换机,而非直接发送到队列。交换机根据绑定规则(Binding Key)和路由键将消息转发到对应的队列。这样,生产者无需关心消息的具体消费队列,消费者也无需了解消息的来源。

    交换机的类型

    交换机有四种主要类型,每种类型适用于不同的场景:

    1. Direct直连交换机

    就如同他的名字,直接将内容放入到指定的queue队列。消息通过完全匹配路由键转发到绑定的队列。适用于点对点精确路由场景,例如订单系统根据订单ID分发消息。

    String queueName = "hello.queue1";
    rabbitTemplate.convertAndSend(queueName,"hello,everyone");

    或者指定routingKey来发送

    String exchangeName = "hello.direct";
    rabbitTemplate.convertAndSend(exchangeName ,"nihao,everyone");

    通过下图可以看到,我指定hello.queue1的routingkey为nihao,所以只会发送给hello.queue1。 是可以多个queue绑定相同的routingKey的 在这里插入图片描述

    2. Fanout广播交换机

    广播,顾名思义就是将绑定的queue都发。将消息广播到所有绑定的队列。适用于发布/订阅模式,例如系统日志广播或实时通知。

    String exchangeName = "hello.fanout";
    rabbitTemplate.convertAndSend(exchangeName ,null,"hello,everyone");

    Consumer

    package com.itheima.consumer.listener;

    import org.apache.logging.log4j.LogManager;
    import org.apache.logging.log4j.Logger;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;

    @Component
    public class RabbitMQListener {
    @RabbitListener(queues = "hello.queue1")
    public void onMessage(String message) {
    Logger logger = LogManager.getLogger(RabbitMQListener.class);
    logger.info("hello.queue1获取到的信息为{}",message);
    }
    @RabbitListener(queues = "hello.queue2")
    public void onMessage2(String message) {
    Logger logger = LogManager.getLogger(RabbitMQListener.class);
    logger.info("hello.queue2获取到的信息为{}",message);
    }
    }

    3. topic主题交换机

    如果说上面的direct交换机是选择队列去发送,那这个topic交换机就是他的升级版,topic交换机可以将routingKey设置为通配符样式。

    发送到类型是 topic 交换机的消息的routing key 不能随意写,必须满足一定的要求,它必须是一个单词列表,并且以点号. 分隔开 。

    这些单词可以是任意单词,比如说:“stock.usd.nyse”,“nyse.vmw”,“quick.orange.rabbit” 这种类型的。当然这个单词列表最多不能超过 255 个字节。

    在这个规则列表中,其中有两个替换符是特别需要注意的:

  • *(星号) 可以代替一个单词,注意是一个单词,不是一个字母
  • #(井号) 可以替代零个或多个单词,注意是一个单词,不是一个字母
  • 在这里插入图片描述

    下面代码就可以发送给这两个queue。

    package com.itheima.publisher;

    import org.junit.jupiter.api.Test;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.boot.test.context.SpringBootTest;

    import javax.annotation.Resource;

    @SpringBootTest
    public class PublisherTest {
    @Resource
    private RabbitTemplate rabbitTemplate;
    @Test
    public void publisherTest(){
    String exchangeName = "hello.topic";
    rabbitTemplate.convertAndSend(exchangeName ,"china.news","hello,everyone");
    }
    }

    声明队列和交换机

    如果每次都在rabbit的控制台去创建队列交换机啥的,未免太麻烦了,spring-amqp为我们封装了代码创建的方法。

    package com.itheima.consumer.config;

    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.FanoutExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Configuration
    public class RabbitConfig {
    @Bean
    public Queue helloQueue1(){
    return new Queue("hello.queue1");
    }

    @Bean
    public Queue helloQueue2(){
    return new Queue("hello.queue2");
    }

    @Bean
    public FanoutExchange helloFanout(){
    return new FanoutExchange("hello.fanout");
    }
    @Bean
    public Binding helloBinding1(){
    return BindingBuilder.bind(helloQueue1()).to(helloFanout());
    }
    @Bean
    public Binding helloBinding2(){
    return BindingBuilder.bind(helloQueue2()).to(helloFanout());
    }

    @Bean
    public DirectExchange helloDirect(){
    return new DirectExchange("hello.direct");
    }
    }

    通过这个操作,就可以实现让服务去创建这些内容。

    注解方式

    amqp为我们提供了更方便的方式,注解方式。 通过下列的操作,可以直接在监听类上完成queue和exchange的创建和绑定。 在这里插入图片描述

    queue就是监听的队列,exchange就是队列要绑定的交换机,key就是routingKey。

    package com.itheima.consumer.listener;

    import org.apache.logging.log4j.LogManager;
    import org.apache.logging.log4j.Logger;
    import org.springframework.amqp.core.ExchangeTypes;
    import org.springframework.amqp.rabbit.annotation.Exchange;
    import org.springframework.amqp.rabbit.annotation.Queue;
    import org.springframework.amqp.rabbit.annotation.QueueBinding;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;

    @Component
    public class RabbitMQListener {
    @RabbitListener(bindings = @QueueBinding(
    value = @Queue(name = "hello.queue1",durable = "true"),
    exchange = @Exchange(name = "hello.direct",type = ExchangeTypes.DIRECT),
    key = {"nihao","zaijian"}
    ))
    public void onMessage(String message) {
    Logger logger = LogManager.getLogger(RabbitMQListener.class);
    logger.info("hello.queue1获取到的信息为{}",message);
    }
    @RabbitListener(bindings = @QueueBinding(
    value = @Queue(name = "hello.queue1",durable = "true"),
    exchange = @Exchange(name = "hello.direct",type = ExchangeTypes.DIRECT),
    key = {"china","emilia"}
    ))
    public void onMessage2(String message) {
    Logger logger = LogManager.getLogger(RabbitMQListener.class);
    logger.info("hello.queue2获取到的信息为{}",message);
    }

    @RabbitListener(bindings = @QueueBinding(
    value = @Queue(name = "hello.queue3",durable = "true"),
    exchange = @Exchange(name = "hello.topic",type = ExchangeTypes.TOPIC),
    key = {"china.#","#.emilia"}
    ))
    public void onMessage3(String message) {
    Logger logger = LogManager.getLogger(RabbitMQListener.class);
    logger.info("hello.queue3获取到的信息为{}",message);
    }
    }

    消息转换器

    rabbitMQ自带的消息转换器是臃肿的。

    我通过下面的代码向交换机发送内容。

    package com.itheima.publisher;

    import org.junit.jupiter.api.Test;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.boot.test.context.SpringBootTest;

    import javax.annotation.Resource;
    import java.util.HashMap;

    @SpringBootTest
    public class PublisherTest {
    @Resource
    private RabbitTemplate rabbitTemplate;
    @Test
    public void publisherTest(){
    HashMap<String, Object> map = new HashMap<>();
    map.put("id","1");
    map.put("name","cmc");
    map.put("sex","男");
    String exchangeName = "hello.topic";
    rabbitTemplate.convertAndSend(exchangeName ,"i love u.emilia",map);
    }
    }

    这个是交换机里面的结果,可以看见第二个结果是我使用了jackson的消息转换,将内容转换为了json,明显大小变小了很多,所以更加推荐使用json消息转换。 在这里插入图片描述

    使用下面的代码,将amqp的消息转换实现注册为bean,amqp就会自动转换消息接受类型了。

    package com.itheima.consumer.config;

    import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
    import org.springframework.amqp.support.converter.MessageConverter;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Configuration
    public class RabbitConfig {
    @Bean
    public MessageConverter messageConverter(){
    return new Jackson2JsonMessageConverter();
    }
    }

    可靠性

    我们知道,如果rabbit没有连接上,或者中途中断,我们要保证数据可靠性问题。 amqp为我们在配置yml中提供了网络断开重连配置。 在这里插入图片描述

    logging:
    pattern:
    dateformat: MMdd HH:mm:ss:SSS
    spring:
    rabbitmq:
    host: localhost
    port: 5672
    virtual-host: /cmc
    username: cmc
    password: 123
    connection-timeout: 1s
    template:
    retry:
    enabled: true
    multiplier: 2 #倍数时间重试

    不过要注意

    当网络不稳定的时候,利用重试机制可以有效提高消息发送的成功率。不过SpringAMOP提供的重试机制是阻塞式的重试,也就是说多次重试等待的过程中,当前线程是被阻塞的,会影响业务性能。 如果对于业务性能有要求,建议禁用重试机制。如果一定要使用,请合理配置等待时长和重试次数,当然也可以考虑使用异步线程来执行发送消息的代码。

    生产者确认

    发送一个消息后,交换机会给生产者返回响应,可以以此来判断是否成功。 在这里插入图片描述

    spring:
    rabbitmq:
    host: localhost
    port: 5672
    virtual-host: /cmc
    username: cmc
    password: 123
    connection-timeout: 1s
    template:
    retry:
    enabled: true
    multiplier: 2
    publisher-confirm-type: correlated
    publisher-returns: true

    他是基于回调机制,所以需要编写回调callback方法。

    这个是关注路由失败的,当内容进入到交换机但是路由到队列失败的时候,会执行这个回调方法,然后发送者确认会返回ack。

    package com.itheima.publisher.config;

    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.core.ReturnedMessage;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.beans.BeansException;
    import org.springframework.context.ApplicationContext;
    import org.springframework.context.ApplicationContextAware;
    import org.springframework.context.annotation.Configuration;

    @Slf4j
    @Configuration
    public class MQConfig implements ApplicationContextAware {

    @Override
    public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
    RabbitTemplate rabbitTemplate = applicationContext.getBean(RabbitTemplate.class);
    //配置回调
    rabbitTemplate.setReturnsCallback(returnedMessage -> log.info("收到消息的return callback,exchange:{},key:{},msg:{}," +
    "code:{},text:{}.",
    returnedMessage.getExchange(),
    returnedMessage.getRoutingKey(),
    returnedMessage.getMessage(),
    returnedMessage.getReplyCode(),
    returnedMessage.getReplyText()));
    }
    }

    我们可以从下面看到,我们加入了 CorrelationData这个对象,这个对象是为我们提供发送者确认回调的,只要内容发送出去就会返回内容。

    当内容发送到达exchange,但是没有路由成功,他会返回ack,由上面的配置类返回错误。 如果是都到达了直接返回ack。 其他原因都会发送nack。

    package com.itheima.publisher;

    import lombok.extern.slf4j.Slf4j;
    import org.junit.jupiter.api.Test;
    import org.springframework.amqp.rabbit.connection.CorrelationData;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.boot.test.context.SpringBootTest;
    import org.springframework.util.concurrent.ListenableFutureCallback;

    import javax.annotation.Resource;
    import java.util.HashMap;
    import java.util.UUID;

    @SpringBootTest
    @Slf4j
    public class PublisherTest {
    @Resource
    private RabbitTemplate rabbitTemplate;
    @Test
    public void publisherTest(){
    HashMap<String, Object> map = new HashMap<>();
    map.put("id","1");
    map.put("name","cmc");
    map.put("sex","男");
    String exchangeName = "hello.direct";
    CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
    correlationData.getFuture().addCallback(new ListenableFutureCallback<CorrelationData.Confirm>() {
    @Override
    public void onFailure(Throwable ex) {
    log.error("消息回调失败");
    }

    @Override
    public void onSuccess(CorrelationData.Confirm result) {
    log.info("收到confirm callback回执");
    if (result.isAck()) {
    log.info("消息发送成功,收到ack");
    }else {
    log.error("消息发送失败,收到nack,reason:{}",result.getReason());
    }
    }
    });

    rabbitTemplate.convertAndSend(exchangeName ,"sfsfw",map,correlationData);
    }
    }

    数据持久化

    默认发送到mq的数据是非持久化的。

  • 当mq遇到宕机,会使得保存在内存中的数据丢失。
  • 而且如果数据量大,可能会让mq将数据移动到磁盘,这个过程可能pageOut。
  • 所以提前将数据库持久化是重要的。

    那如何持久化呢?rabbitMQ为我们封装了方法

    package com.itheima.publisher;

    import com.alibaba.fastjson2.JSONObject;
    import lombok.extern.log4j.Log4j2;
    import lombok.extern.slf4j.Slf4j;
    import org.junit.jupiter.api.Test;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.core.MessageBuilder;
    import org.springframework.amqp.core.MessageDeliveryMode;
    import org.springframework.amqp.rabbit.connection.CorrelationData;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.boot.test.context.SpringBootTest;
    import org.springframework.util.concurrent.ListenableFutureCallback;

    import javax.annotation.Resource;
    import java.nio.charset.StandardCharsets;
    import java.util.HashMap;
    import java.util.UUID;

    @SpringBootTest
    @Slf4j
    public class PublisherTest {
    @Resource
    private RabbitTemplate rabbitTemplate;

    @Test
    public void publisherTest() {
    HashMap<String, Object> map = new HashMap<>();
    map.put("id", "1");
    map.put("name", "cmc");
    map.put("sex", "男");
    String exchangeName = "hello.direct";
    CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
    correlationData.getFuture().addCallback(new ListenableFutureCallback<CorrelationData.Confirm>() {
    @Override
    public void onFailure(Throwable ex) {
    log.error("消息回调失败");
    }

    @Override
    public void onSuccess(CorrelationData.Confirm result) {
    log.info("收到confirm callback回执");
    if (result.isAck()) {
    log.info("消息发送成功,收到ack");
    } else {
    log.error("消息发送失败,收到nack,reason:{}", result.getReason());
    }
    }
    });

    //这里就是
    Message message = MessageBuilder.withBody(
    JSONObject.toJSONString(map).getBytes(StandardCharsets.UTF_8))
    .setDeliveryMode(MessageDeliveryMode.PERSISTENT)
    .build();
    rabbitTemplate.convertAndSend(exchangeName, "emilia", message, correlationData);
    Message message1 = MessageBuilder.withBody(
    "cmcnb".getBytes(StandardCharsets.UTF_8))
    .setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT)
    .build();
    rabbitTemplate.convertAndSend(exchangeName, "emilia", message1, correlationData);
    }
    }

    通过将发送的信息转换为mode=2,来让消息是持久化,不过默认发送的消息就是2。 在这里插入图片描述 但是因为他每一次都要往磁盘里面写数据,所以性能不是很好。 在这里插入图片描述

    Lazy Queue

    rabbitMQ在3.6增加了一种新概念,叫惰性队列LazyQueue。 惰性队列的特征如下:

    • 接收到消息后直接存入磁盘而非内存(内存中只保留最近的消息,默认2048条)
    • 消费者要消费消息时才会从磁盘中读取并加载到内存
    • 支持数百万条的消息存储

    在3.12版本后,所有队列都是Lazy Queue模式,无法更改

    在这里插入图片描述

    package com.itheima.consumer.config;

    import org.springframework.amqp.core.Queue;
    import org.springframework.amqp.core.QueueBuilder;
    import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
    import org.springframework.amqp.support.converter.MessageConverter;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Configuration
    public class RabbitConfig {
    @Bean
    public MessageConverter messageConverter() {
    return new Jackson2JsonMessageConverter();
    }

    @Bean
    public Queue queue() {
    return QueueBuilder.durable("hello.queue1")
    .lazy().
    build();
    }
    }

    或者listener监听直接

    package com.itheima.consumer.listener;

    import org.apache.logging.log4j.LogManager;
    import org.apache.logging.log4j.Logger;
    import org.springframework.amqp.core.ExchangeTypes;
    import org.springframework.amqp.rabbit.annotation.*;
    import org.springframework.stereotype.Component;

    import java.util.Map;

    @Component
    public class RabbitMQListener {

    @RabbitListener(bindings = @QueueBinding(
    value = @Queue(name = "hello.queue2",durable = "true"),
    exchange = @Exchange(name = "hello.direct"),
    arguments = @Argument(name = "x-queue-mode",value = "lazy"),
    key = {"china","emilia"}
    ))
    public void onMessage2(String message) {
    Logger logger = LogManager.getLogger(RabbitMQListener.class);
    logger.info("hello.queue2获取到的信息为{}",message);
    }
    }

    消费者确认机制

    为了确认消费者是否成功处理消息,RabbitMQ提供了消费者确认机制(ConsumerAcknowledgement)。当消费者处理消息结束后,应该向RabbitMO发送一个回执,告知RabbitM0自己消息处理状态。

    回执有三种可选值:

    • ack:成功处理消息,RabbitMQ从队列中删除该消息
    • nack:消息处理失败,RabbitMQ需要再次投递消息
    • reject:消息处理失败并拒绝该消息,RabbitMQ从队列中删除该消息 在这里插入图片描述

    SpringAMQP已经实现了消息确认功能。并允许我们通过配置文件选择ACK处理方式, 有三种方式

    • none:不处理。即消息投递给消费者后立刻ack,消息会立刻从MQ删除。非常不安全,不建议使用
    • manual:手动模式。需要自己在业务代码中调用api,发送ack或reject,存在业务入侵,但更灵活
    • auto:自动模式。SpringAMQP利用AOP对我们的消息处理逻辑做了环绕增强,当业务正常执行时则自动返回ack.当业务出现异常时,根据异常判断返回不同。 在这里插入图片描述

    结果: 如果是业务异常,会自动返回nack 如果是消息处理或校验异常,自动返回

    消费者失败重试机制

    当消费者出现异常后,消息会不断requeue(重新入队)到队列,再重新发送给消费者,然后再次异常,再次requeue无限循环,导致mq的消息处理飙升,带来不必要的压力。我们可以利用Spring的retry机制,在消费者出现异常时利用本地重试,而不是无限制的requeue到mq队列:

    logging:
    pattern:
    dateformat: MMdd HH:mm:ss:SSS
    spring:
    rabbitmq:
    host: localhost
    port: 5672
    virtualhost: /cmc
    username: cmc
    password: 123
    listener:
    simple:
    retry:
    enabled: true #打开重试
    multiplier: 2 #下次失败的等待时间倍数
    initialinterval: 1000ms #初始失败等待时间
    maxattempts: 3 #最大尝试次数

    在这里插入图片描述

    消费者失败消息处理策略

    在开启重试模式后,重试次数耗尽,如果消息依然失败,则需要有MessageRecoverer接口来处理,它包含三种不同的实现:

    • RejectAndDontRequeueRecoverer:重试耗尽后,直接reiect,丢弃消息。默认就是这种方式。
    • ImmediateRequeueMessageRecoverer:重试耗尽后,返回nack,消息重新入队
    • RepublishMessageRecoverer:重试耗尽后,将失败消息投递到指定的交换机

    这里就实现一下RepublishMessageRecoverer: 这个意思是创建一个error队列来存储获取的错误消息

    package com.itheima.consumer.config;
    import org.springframework.amqp.core.*;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.amqp.rabbit.retry.MessageRecoverer;
    import org.springframework.amqp.rabbit.retry.RepublishMessageRecoverer;
    import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
    import org.springframework.amqp.support.converter.MessageConverter;
    import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Configuration
    @ConditionalOnProperty(prefix = "spring.rabbitmq.listener.simple.retry",name = "enabled",havingValue = "true")
    public class RabbitErrorConfig {

    //声明异常交换机
    @Bean
    public DirectExchange errorExchange() {
    return ExchangeBuilder.directExchange("error.direct").durable(true)
    .build();
    }
    //声明异常队列
    @Bean
    public Queue errorQueue() {
    return QueueBuilder.durable("error.queue")
    .lazy().build();
    }
    @Bean
    public Binding errorBinding() {
    return BindingBuilder.bind(errorQueue()).to(errorExchange()).with("");
    }
    //失败消息处理策略
    @Bean
    public MessageRecoverer messageRecoverer(RabbitTemplate rabbitTemplate) {
    return new RepublishMessageRecoverer(rabbitTemplate,"error.direct","");
    }
    }

    消费者如何保证消息一定被消费: 开启消费者机制为auto,由spring来确定消费获取成功为ack,失败为nack。 开启消费者重试机制,当达到最大重试次数后,将内容交给异常交换机,交给人为处理。

    业务幂等性

    幂等是一个数学概念,用函数表达式来描述是这样的:f(x)=f(f(x))。在程序开发中,则是指同一个业务,执行一次或多次对业务状态的影响是一致的。 在这里插入图片描述

    业务唯一ID

    我们一般给mq发送内容,会让他携带一个唯一ID,以此来判断他的唯一性。 我们可以直接封装Jackson2JsonMessageConverter。

    package com.itheima.publisher.config;

    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
    import org.springframework.amqp.support.converter.MessageConverter;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Slf4j
    @Configuration
    public class MQConfig {

    @Bean
    public MessageConverter messageConverter() {
    Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter();
    //打开messageID
    converter.setCreateMessageIds(true);
    return converter;
    }

    }

    但如果你在发送内容的时候是自己封装了message,像这个样子。 在这里插入图片描述 那么你就要自己在message上面加上设置messageID,因为这样会跳过amqp的jacksonMessage封装。

    延迟消息

    延迟消息是一种特殊的消息类型,允许消息在指定的时间后才被消费者处理。这种机制在现代分布式系统中非常重要,尤其是在需要定时任务、重试机制或缓冲处理的场景中。

    延迟消息的核心思想是:生产者发送消息时指定一个延迟时间,消息会在队列中暂存,直到延迟时间到期后才被消费者消费。

    死信交换机

    一个队列中的消息满足下列情况之一时,就会成为死信(dead letter): 消费者使用basic.reject或 basic.nack声明消费失败,并且消息的requeue参数设置为false消息是一个过期消息(达到了队列或消息本身设置的过期时间),超时无人消费。 要投递的队列消息堆积满了,最早的消息可能成为死信 如果队列通过dead-letter-exchange属性指定了一个交换机,那么该队列中的死信就会投递到这个交换机中。这个交换机称为死信交换机(Dead LetterExchange,简称DLX)。 在这里插入图片描述

    延迟插件

    rabbitMQ为我们提供了延迟插件。

    下载并安装延迟插件

    RabbitMQ 的官方推出了一个插件,原生支持延迟消息功能。该插件的原理是设计了一种支持延迟消息功能的交换机,当消息投递到交换机后,可以将消息暂存一段时间,时间到了之后再将消息投递到队列中

    插件的下载地址:rabbitmq-delayed-message-exchange

    在这里插入图片描述

    下载完插件后,将这个插件放到到 RabbitMQ 的插件的安装目录,我的是windows目录就是在下图这里。 在这里插入图片描述

    这样你就可以在创建交换机的时候看见创建延迟交换机的类型。 在这里插入图片描述

    使用

    通过consumer来进行监听并且创建 在这里插入图片描述 publisher来进行发送 在这里插入图片描述

    package com.itheima.publisher;

    import com.alibaba.fastjson2.JSONObject;
    import lombok.extern.log4j.Log4j2;
    import lombok.extern.slf4j.Slf4j;
    import org.junit.jupiter.api.Test;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.core.MessageBuilder;
    import org.springframework.amqp.core.MessageDeliveryMode;
    import org.springframework.amqp.rabbit.connection.CorrelationData;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
    import org.springframework.boot.test.context.SpringBootTest;
    import org.springframework.util.concurrent.ListenableFutureCallback;

    import javax.annotation.Resource;
    import java.nio.charset.StandardCharsets;
    import java.util.HashMap;
    import java.util.UUID;

    @SpringBootTest
    @Slf4j
    public class PublisherTest {
    @Resource
    private RabbitTemplate rabbitTemplate;

    @Test
    public void publisherTest() throws InterruptedException {
    CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
    correlationData.getFuture().addCallback(new ListenableFutureCallback<CorrelationData.Confirm>() {
    @Override
    public void onFailure(Throwable ex) {
    log.error("消息回调失败");
    }

    @Override
    public void onSuccess(CorrelationData.Confirm result) {
    log.info("收到confirm callback回执");
    if (result.isAck()) {
    log.info("消息发送成功,收到ack");
    } else {
    log.error("消息发送失败,收到nack,reason:{}", result.getReason());
    }
    }
    });
    Message message1 = MessageBuilder.withBody(
    "cmcnb".getBytes(StandardCharsets.UTF_8))
    .setMessageId(UUID.randomUUID().toString())
    .setExpiration("3000")
    .build();
    rabbitTemplate.convertAndSend(exchangeName, "huawei", message1, correlationData);
    }
    }

    也可以以这种方式来封装消息头 在这里插入图片描述

    package com.itheima.publisher;

    import com.alibaba.fastjson2.JSONObject;
    import lombok.extern.log4j.Log4j2;
    import lombok.extern.slf4j.Slf4j;
    import org.junit.jupiter.api.Test;
    import org.springframework.amqp.AmqpException;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.core.MessageBuilder;
    import org.springframework.amqp.core.MessageDeliveryMode;
    import org.springframework.amqp.core.MessagePostProcessor;
    import org.springframework.amqp.rabbit.connection.CorrelationData;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
    import org.springframework.boot.test.context.SpringBootTest;
    import org.springframework.util.concurrent.ListenableFutureCallback;

    import javax.annotation.Resource;
    import java.nio.charset.StandardCharsets;
    import java.util.HashMap;
    import java.util.UUID;

    @SpringBootTest
    @Slf4j
    public class PublisherTest {
    @Resource
    private RabbitTemplate rabbitTemplate;

    @Test
    public void publisherTest() throws InterruptedException {
    HashMap<String, Object> map = new HashMap<>();
    map.put("id", "1");
    map.put("name", "cmc");
    map.put("sex", "男");
    String exchangeName = "dealy.direct";
    CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
    correlationData.getFuture().addCallback(new ListenableFutureCallback<CorrelationData.Confirm>() {
    @Override
    public void onFailure(Throwable ex) {
    log.error("消息回调失败");
    }

    @Override
    public void onSuccess(CorrelationData.Confirm result) {
    log.info("收到confirm callback回执");
    if (result.isAck()) {
    log.info("消息发送成功,收到ack");
    } else {
    log.error("消息发送失败,收到nack,reason:{}", result.getReason());
    }
    }
    });

    Message message1 = MessageBuilder.withBody(
    "cmcnb".getBytes(StandardCharsets.UTF_8))
    .build();
    rabbitTemplate.convertAndSend(exchangeName, "huawei", message1, new MessagePostProcessor() {
    @Override
    public Message postProcessMessage(Message message) throws AmqpException {
    message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);//持久化
    message.getMessageProperties().setContentEncoding("UTF-8");
    message.getMessageProperties().setDelay(3000);//延迟时间
    message.getMessageProperties().setMessageId(UUID.randomUUID().toString());//messageID
    return message;
    }
    }, correlationData);
    }
    }

    延迟消息的原理和缺点

    RabbitMQ 的延迟消息是怎么实现的呢?RabbitMQ 会自动维护一个时钟,这个时钟每隔一秒就跳动一次,如果对时钟的精度要求比较高的,可能还要精确到毫秒,甚至纳秒

    RabbitMQ 会为发送到交换机的每一条延迟消息创建一个时钟,时钟运行的过程中需要 CPU 不断地进行计算。发送到交换机的延迟消息数越多,RabbitMQ 需要维护的时钟就越多,对 CPU 的占用率就越高(Spring 提供的定时任务的原理也是类似)

    定时任务属于 CPU 密集型任务,中间涉及到的计算过程对 CPU 来说压力是很大的,所以说,采用延迟消息会给服务器的 CPU 带来更大的压力。当交换机中有非常多的延迟消息时,对 CPU 的压力就会特别大

    所以说,延迟消息适用于延迟时间较短的场景

    实景操作 取消30分钟订单

    设置 30 分钟后检测订单支付状态实现起来非常简单,但是存在两个问题:

  • 如果并发较高,30分钟可能堆积消息过多,对 MQ 压力很大
  • 大多数订单在下单后 1 分钟内就会支付,但消息需要在 MQ,中等待30分钟,浪费资源。
  • 通过将10分钟裁成不同的分片,在每一次分片都去查询订单状态是否完成,如果没完成就执行下一个分片时间,直到delayTIme用完,取消订单,这样可以在订单完成的时候直接取消MQ的延迟信息,不占用MQ空间。 在这里插入图片描述 下面是封装的分片延迟信息Message

    import java.util.ArrayList;
    import java.util.Arrays;
    import java.util.List;

    public class MultipleDelayMessage<T> {

    private T data;

    private List<Long> delayMillis;

    public MultipleDelayMessage() {

    }

    public MultipleDelayMessage(T data, Long... delayMillis) {
    this.data = data;
    this.delayMillis = new ArrayList<>(Arrays.asList(delayMillis));
    }

    public MultipleDelayMessage(T data, List<Long> delayMillis) {
    this.data = data;
    this.delayMillis = delayMillis;
    }

    public static <T> MultipleDelayMessage<T> of(T data, Long... delayMillis) {
    return new MultipleDelayMessage<>(data, new ArrayList<>(Arrays.asList(delayMillis)));
    }

    public static <T> MultipleDelayMessage<T> of(T data, List<Long> delayMillis) {
    return new MultipleDelayMessage<>(data, delayMillis);
    }

    public boolean hasNextDelay() {
    return !delayMillis.isEmpty();
    }

    public Long removeNextDelay() {
    return delayMillis.remove(0);
    }

    public T getData() {
    return data;
    }

    public void setData(T data) {
    this.data = data;
    }

    public List<Long> getDelayMillis() {
    return delayMillis;
    }

    public void setDelayMillis(List<Long> delayMillis) {
    this.delayMillis = delayMillis;
    }

    @Override
    public String toString() {
    return "MultipleDelayMessage{" +
    "data=" + data +
    ", delayMillis=" + delayMillis +
    '}';
    }

    }

    流程图在下面,就如我上面所说。 在这里插入图片描述

    赞(0)
    未经允许不得转载:171主机测评 » RabbitMQ是什么?如何使用
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址