zoukankan      html  css  js  c++  java
  • SpringBoot整合RabbitMQ实现延时队列

    SpringBoot之RabbitMQ实现两种延时队列

    参考:https://blog.csdn.net/qq_37892957/article/details/89296157
    参考:https://blog.csdn.net/u014308482/article/details/53036770

    原理

    Rabbit本身不支持延时队列。但是,我们可以利用其死信(dead letter)投递机制来实现延时效果。
    当队列中的消息未被正常消费时,就会被标记为死信(dead letter)。

    概念

    即消息放入队列之后,不立即进行消费,而是过了一段时间之后才进行消费。

    场景

    订单场景中就会用到延时队列,比如某个用户在15分钟内没有去完成订单,这个时候我们就需要去自动的去取消这边订单。

    实现

    1. 方法一:
      利用到RabbitMQ的两个特性:
      • Time To Live (TTL)
      • Dead Letter Exchanges (DLX)
    2. 方法二:
      利用RabbitMQ中的插件x-delay-message

    TTL

    RabbitMQ可以针对Queue设置x-expires 或者 针对Message设置 x-message-ttl,来控制消息的生存时间,如果超时(两者同时设置以最先到期的时间为准),则消息变为死信(dead letter)。

    RabbitMQ针对队列中的消息过期时间有两种方法可以设置。

    • 通过队列属性设置,队列中所有消息都有相同的过期时间。
    • 对消息进行单独设置,每条消息TTL可以不同。

    DLX

    RabbitMQ的Queue可以配置 x-dead-letter-exchangex-dead-letter-routing-key两个(可选)参数,如果队列内出现了死信(dead letter),则按照这两个参数重新路由转发到指定的队列。

    • x-dead-letter-exchange:出现dead letter之后将dead letter重新发送到指定exchange
    • x-dead-letter-routing-key:出现dead letter之后将dead letter重新按照指定的routing-key发送

    如果同时使用,则消息的过期时间以两者之间TTL较小的那个数值为准。消息在队列的生存时间一旦超过设置的TTL值,就成为死信(dead letter)。

    注意:不建议使用TTL这种方式来实现延迟队列,过期的消息变为死信,进入死信接收队列,而这个队列就是普通的队列,如果这个队列拥塞了很多死信,那么死信出队列的顺序就是其进入死信接收队列的顺序。

    TTL & DLX

    1. RabbitMQ的Maven依赖
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
    
    1. RabbitMQ在application.yml文件中的配置
    spring:
        rabbitmq:
          # rabbitmq的服务地址
          host: 192.168.1.160
          port: 5672
          username: guest
          password: guest
    
    1. 队列配置
    import org.springframework.amqp.core.*;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.boot.autoconfigure.amqp.SimpleRabbitListenerContainerFactoryConfigurer;
    import org.springframework.context.annotation.Bean;
     
    import java.util.HashMap;
    import java.util.Map;
    
    @Configuration
    public class RabbitmqConfig {
       /**
         * 死信交换机
         */
        @Bean
        public DirectExchange userOrderDelayExchange() {
            return new DirectExchange("user.order.delay_exchange", true, false);
        }
     
        /**
         * 死信队列
         */
        @Bean
        public Queue userOrderDelayQueue() {
            Map<String, Object> map = new HashMap<>(3);
            // 设置15分钟过期时间
            map.put("x-message-ttl", 900000);
            map.put("x-dead-letter-exchange", "user.order.receive_exchange");
            map.put("x-dead-letter-routing-key", "user.order.receive_key");
            return new Queue("user.order.delay_queue", true, false, false, map);
        }
     
        /**
         * 给死信队列绑定交换机
         */
        @Bean
        public Binding userOrderDelayBinding() {
            return BindingBuilder.bind(userOrderDelayQueue()).to(userOrderDelayExchange()).with("user.order.delay_key");
        }
     
        /**
         * 死信接收交换机
         */
        @Bean
        public DirectExchange userOrderReceiveExchange() {
            return new DirectExchange("user.order.receive_exchange", true, false);
        }
     
        /**
         * 死信接收队列,用于接收死信,该队列为正常队列,进入该队列的消息会被立即消费
         */
        @Bean
        public Queue userOrderReceiveQueue() {
            return new Queue("user.order.receive_queue");
        }
     
        /**
         * 给死信交换机绑定消费队列
         */
        @Bean
        public Binding userOrderReceiveBinding() {
            return BindingBuilder.bind(userOrderReceiveQueue()).to(userOrderReceiveExchange()).with("user.order.receive_key");
        }
     
    }
    
    1. 发送端代码(Sender)
        @Autowired
        private RabbitTemplate rabbitTemplate;
        @GetMapping("test")
        public String createOrderTest() {
            OrderMaster orderMaster = new OrderMaster();
            // 未支付
            orderMaster.setOrderStatus(0);
            // 未支付
            orderMaster.setPayStatus(0);
            orderMaster.setBuyerName("张三");
            orderMaster.setBuyerAddress("湖南长沙");
            orderMaster.setBuyerPhone("186981578424");
            orderMaster.setOrderAmount(BigDecimal.ZERO);
            orderMaster.setCreateTime(DateUtils.getCurrentDate());
            orderMaster.setOrderId(UUID.randomUUID().toString().replaceAll("-", ""));
            orderMasterService.insert(orderMaster);
            // TODO:设置超时,用mq处理已超时的下单记录(一旦记录超时,则处理为无效)
            rabbitTemplate.convertAndSend("user.order.delay_exchange", "user.order.delay_key", orderMaster, message -> {
                message.getMessageProperties().setExpiration("300000");
                return message;
            });
            return "创建订单成功";
        }
    
    1. 消息接收端(Receiver)
    import com.bean.springcloudcommon.model.OrderMaster;
    import com.rabbitmq.client.Channel;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.stereotype.Component;
     
    import java.io.IOException;
    import java.util.Objects;
    
    @Component
    public class OrderReceiver {
        @Autowired
        private OrderMasterService orderMasterService;
     
        // 监听消息队列
        @RabbitListener(queues = "user.order.receive_queue")
        public void consumeMessage(OrderMaster order) throws IOException {
            try {
                // 如果订单状态不是0 说明订单已经被其他消费队列改动过了 加一个状态用来判断集群状态的情况
                if (Objects.equals(0,order.getOrderStatus())) {
                    // 设置订单过去状态
                    order.setOrderStatus(-1);
                    System.out.println(order.getBuyerName());
                    orderMasterService.updateByPrimaryKeySelective(order);
                }
            } catch (Exception e) {
                e.printStackTrace()
            }
        }
    }
    

    x-delay-message插件(推荐使用)

    1. 插件下载
      在rabbitmq 3.5.7及以上的版本提供了一个插件 rabbitmq-delayed-message-exchange 来实现延迟队列功能。同时插件依赖 Erlang/OPT 18.0 及以上。

    插件源码地址:点击这里

    插件下载地址:点击这里

    1. 插件安装
    # 进入插件安装目录(可以查看一下当前已存在的插件)
    cd {rabbitmq-server}/plugins/
    
    # 下载插件 rabbitmq_delayed_message_exchange
    wget https://bintray.com/rabbitmq/community-plugins/download_file?file_path=rabbitmq_delayed_message_exchange-0.0.1.ez
    
    # 启用插件
    rabbitmq-plugins enable rabbitmq_delayed_message_exchange
    
    # 关闭插件
    rabbitmq-plugins disable rabbitmq_delayed_message_exchange
    
    1. RabbitMQ的Maven依赖
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
    
    1. RabbitMQ在application.yml文件中的配置
    spring:
        rabbitmq:
          # rabbitmq的服务地址
          host: 192.168.1.160
          port: 5672
          username: guest
          password: guest
    
    1. 配置队列
    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.CustomExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
     
    import java.util.HashMap;
    import java.util.Map;
    
    @Configuration
    public class RabbitmqConfig {
        /**
         * 延时队列交换机
         * 注意这里的交换机类型:CustomExchange
         */
        @Bean
        public CustomExchange delayExchange() {
            Map<String, Object> args = new HashMap<>();
            args.put("x-delayed-type", "direct");
            // 属性参数 交换机名称 交换机类型 是否持久化 是否自动删除 配置参数
            return new CustomExchange("delay_exchange", "x-delayed-message", true, false, args);
        }
    
        /**
         * 延时队列
         */
        @Bean
        public Queue delayQueue() {
            // 属性参数 队列名称 是否持久化
            return new Queue("delay_queue", true);
        }
    
        /**
         * 给延时队列绑定交换机
         */
        @Bean
        public Binding cfgDelayBinding() {
            return BindingBuilder.bind(delayQueue()).to(delayExchange()).with("delay_key").noargs();
        }
    }
    
    1. 消息发送端(Sender)
    @GetMapping("test/{time}/{name}")
    public String createOrderTest(@PathVariable("time") Integer time, @PathVariable("name") String name) {
        OrderMaster orderMaster = new OrderMaster();
        // 订单未完成
        orderMaster.setOrderStatus(0);
        // 未付款
        orderMaster.setPayStatus(0);
        orderMaster.setBuyerName(name);
        orderMaster.setBuyerAddress("湖南长沙");
        orderMaster.setBuyerPhone("手机号");
        orderMaster.setOrderAmount(BigDecimal.ZERO);
        orderMaster.setCreateTime(DateUtils.getCurrentDate());
        orderMaster.setOrderId(UUID.randomUUID().toString().replaceAll("-", ""));
        orderMasterService.insert(orderMaster);
         // 第一个参数是前面RabbitMqConfig的交换机名称 第二个参数的路由名称 第三个参数是传递的参数 第四个参数是配置属性
        this.rabbitTemplate.convertAndSend(
                "delay_exchange",
                "delay_key",
                orderMaster,
                message -> {
                    // 配置消息的过期时间
                    message.getMessageProperties().setDelay(time);
                    return message;
                }
        );
        return "创建订单成功";
    }
    
    1. 消息接收端(Receiver)
    import com.bean.springcloudcommon.model.OrderMaster;
    import com.rabbitmq.client.Channel;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.stereotype.Component;
     
    import java.io.IOException;
    import java.util.Objects;
    
    @Component
    public class OrderReceiver {
        @Autowired
        private OrderMasterService orderMasterService;
     
        // 监听消息队列
        @RabbitListener(queues = "delay_queue")
        public void consumeMessage(OrderMaster order) throws IOException {
            try {
                // 如果订单状态不是0 说明订单已经被其他消费队列改动过了 加一个状态用来判断集群状态的情况
                if (Objects.equals(0,order.getOrderStatus())) {
                    // 设置订单过去状态
                    order.setOrderStatus(-1);
                    System.out.println(order.getBuyerName());
                    orderMasterService.updateByPrimaryKeySelective(order);
                }
            } catch (Exception e) {
                e.printStackTrace()
            }
        }
    }
    
  • 相关阅读:
    postgresql批量删除表
    Oracle迁移至PostgreSQL工具之Ora2Pg
    postgresql获取表最后更新时间(通过发布订阅机制将消息发送给应用程序)
    postgresql获取表最后更新时间(通过表磁盘存储文件时间)
    postgresql获取表最后更新时间(通过触发器将时间写入另外一张表)
    postgresql源码编译安装(centos)
    Java 学习笔记(7)——接口与多态
    Java 学习笔记(6)——继承
    Java 学习笔记(4)——java 常见类
    Java 学习笔记(4)——面向对象
  • 原文地址:https://www.cnblogs.com/jockming/p/13180669.html
Copyright © 2011-2022 走看看