暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

RabbitMQ延迟队列在项目中的实际使用

212
  • 背景介绍

    • 用户在电商系统中会有下单的时候使用优惠券的场景,通常一张优惠券我们只允许被一个商品使用,因此在用户选择某张优惠券的时候我们需要对这个优惠券进行加锁,加锁以后表示这张优惠券后期不能被再次使用了,但是最后有可能用户最后下单失败了,如果用户下单失败了,那么我们应该释放被锁住的优惠券[这样下次用户再次下单就还可以使用这张优惠券了]。

  • 解决方案

    • 用户下单的时候我们需要发送一个优惠券+订单信息的延迟消息,后续通过消费者消费延迟消息的时候去查一下这个订单是否真的付款了,如果没有付款就释放优惠券,如果付款了,就更新优惠卷的状态为已使用。

  • 总体架构

  • 生产者端的业务逻辑【以下三步要求原子操作】

    • 标记该用户的相应的优惠券为已使用状态

    • 根据该用户已经使用的优惠劵和订单信息构建task记录插入并将每一条记录标记为locked 【等待后续消费者端进行验证】

    • 构造延迟消息包括优惠劵id和订单号发送给rabbitmq的延迟队列中。

  • 向springboot注册rabbitMQ的基本信息

    • pom文件中引入依赖,rabbitmq的amqp依赖

    <!--引入AMQP-->
    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
      • 编写rabbitmq的配置文件

      服务器端开启rabbitmq,需要告诉这个微服务我们的rabbitmq的ip和端口
      #消息队列
      rabbitmq:
      host: 8.129.113.233
      port: 5672
      virtual-host: /
      password: password
      username: admin
      #开启手动确认消息
      listener:
      simple:
      acknowledge-mode: manual
        • 在application.yml中配置队列+交换机+路由key的名称(我们会默认路由key和绑定key是一样的)

        #自定义消息队列配置,发送锁定库存消息-》延迟exchange-lock.queue-》死信exchange-release.queue
        mqconfig:
        #延迟队列,不能被监听消费 【普通队列】
        coupon_release_delay_queue: coupon.release.delay.queue


        #延迟队列的消息过期后转发的队列 【死信队列】
        coupon_release_queue: coupon.release.queue


        #交换机
        coupon_event_exchange: coupon.event.exchange


        #进入延迟队列的路由key 【到普通队列的路由key
        coupon_release_delay_routing_key: coupon.release.delay.routing.key


        #消息过期,进入释放死信队列的key 【到私信队列的路由key
        coupon_release_routing_key: coupon.release.routing.key


        #消息过期时间,毫秒,测试改为15秒,用于在普通队列设置参数
        ttl: 15000
        • 编写配置类RabbitMQConfig读取配置文件,交由Spring扫描

          • 延迟队列和普通队列的区别在于需要额外指定一些延迟队列特有的属性【过期时间,消息过期以后需要重新更新的key,死信交换机】

          • 当我们发送消息的时候,自动会在rabbitmq服务端创建好消息和交换机。

          @Configuration
          @Data
          public class RabbitMQConfig {




          /**
          * 交换机
          */
          @Value("${mqconfig.coupon_event_exchange}")
          private String eventExchange;




          /**
          * 第一个队列 延迟队列,即普通队列,
          */
          @Value("${mqconfig.coupon_release_delay_queue}")
          private String couponReleaseDelayQueue;


          /**
          * 第一个队列的路由key
          * 进入队列的路由key
          */
          @Value("${mqconfig.coupon_release_delay_routing_key}")
          private String couponReleaseDelayRoutingKey;




          /**
               * 第二个队列,被监听恢复库存的队列,即死信队列
          */
          @Value("${mqconfig.coupon_release_queue}")
          private String couponReleaseQueue;


          /**
          * 第二个队列的路由key
          *
          * 即进入死信队列的路由key
          */
          @Value("${mqconfig.coupon_release_routing_key}")
          private String couponReleaseRoutingKey;


          /**
          * 过期时间
          */
          @Value("${mqconfig.ttl}")
          private Integer ttl;




          /**
          * 消息转换器
          * @return
          */
          @Bean
          public MessageConverter messageConverter(){
          return new Jackson2JsonMessageConverter();
          }




          /**
          * 创建交换机 Topic类型,也可以用dirct路由
          * 一般一个微服务一个交换机
          * @return
          */
          @Bean
          public Exchange couponEventExchange(){
          return new TopicExchange(eventExchange,true,false);
          }




          /**
          * 延迟队列
          */
          @Bean
          public Queue couponReleaseDelayQueue(){
          /**
          * 延迟队列参数,包括队列中消息过期时间,队列绑定的死信交换机是哪个,队列和死信交换机的路由key是什么
          * 这样在消息过期以后,延迟队列会将过期消息添加路由key投递到死信交换机
          * 这里的args是由队列-》死信交换机
          */
          Map<String,Object> args = new HashMap<>(3);
          args.put("x-message-ttl",ttl);
          args.put("x-dead-letter-routing-key",couponReleaseRoutingKey);
          args.put("x-dead-letter-exchange",eventExchange);


          return new Queue(couponReleaseDelayQueue,true,false,false,args);
          }




          /**
          * 死信队列,普通队列,用于被监听
          */
          @Bean
          public Queue couponReleaseQueue(){
          return new Queue(couponReleaseQueue,true,false,false);
          }




          /**
          * 即延迟队列和交换机的绑定关系建立
          * 这里绑定的是当一个消息到达交换机,交换机根据消息的key投递到哪个队列
          * 这里绑定的关系是由交换机-》队列
          * @return
          */
          @Bean
          public Binding couponReleaseDelayBinding(){


          return new Binding(couponReleaseDelayQueue,Binding.DestinationType.QUEUE,eventExchange,couponReleaseDelayRoutingKey,null);
          }


          /**
          * 死信队列绑定关系建立
          * 里绑定的是当一个消息到达交换机,交换机根据消息的key投递到哪个队列
          * 这里的绑定关系是由交换机-》队列
          * @return
          */
          @Bean
          public Binding couponReleaseBinding(){


          return new Binding(couponReleaseQueue,Binding.DestinationType.QUEUE,eventExchange,couponReleaseRoutingKey,null);
          }


          }
          • 消息类CouponRecordMessage

            • 主要用于封装消息,包括消息id,订单号,优惠券订单号,锁定库存的ID

            • 后续消费者取出消息可以查询去订单数据库查询订单号,检查订单是否下单支付,然后决定锁定的库存记录是否需要回滚

            @Data
            public class CouponRecordMessage {




            /**
            * 消息id
            */
            private String messageId;


            /**
            * 订单号
            */
            private String outTradeNo;




            /**
            * 库存锁定任务id
            */
            private Long taskId;


            }
            • 生产者逻辑

              /**
              * 锁定优惠券
              *
              * 1)锁定优惠券记录
              * 2)task表插入记录
              * 3)发送延迟消息
              *
              * @param recordRequest
              * @return
              */
              @Override
              public JsonData lockCouponRecords(LockCouponRecordRequest recordRequest) {


              LoginUser loginUser = LoginInterceptor.threadLocal.get();


              String orderOutTradeNo = recordRequest.getOrderOutTradeNo();
              List<Long> lockCouponRecordIds = recordRequest.getLockCouponRecordIds();




              int updateRows = couponRecordMapper.lockUseStateBatch(loginUser.getId(),CouponStateEnum.USED.name(),lockCouponRecordIds);


              List<CouponTaskDO> couponTaskDOList = lockCouponRecordIds.stream().map(obj->{
              CouponTaskDO couponTaskDO = new CouponTaskDO();
              couponTaskDO.setCreateTime(new Date());
              couponTaskDO.setOutTradeNo(orderOutTradeNo);
              couponTaskDO.setCouponRecordId(obj);
              couponTaskDO.setLockState(StockTaskStateEnum.LOCK.name());
              return couponTaskDO;
              }).collect(Collectors.toList());


              int insertRows = couponTaskMapper.insertBatch(couponTaskDOList);


              log.info("优惠券记录锁定updateRows={}",updateRows);
              log.info("新增优惠券记录task insertRows={}",insertRows);




              if(lockCouponRecordIds.size() == insertRows && insertRows==updateRows){
              //发送延迟消息


              for(CouponTaskDO couponTaskDO : couponTaskDOList){
              CouponRecordMessage couponRecordMessage = new CouponRecordMessage();
              couponRecordMessage.setOutTradeNo(orderOutTradeNo);
              couponRecordMessage.setTaskId(couponTaskDO.getId());
              rabbitTemplate.convertAndSend(rabbitMQConfig.getEventExchange(),rabbitMQConfig.getCouponReleaseDelayRoutingKey(),couponRecordMessage);
              log.info("优惠券锁定消息发送成功:{}",couponRecordMessage.toString());
              }
              return JsonData.buildSuccess();
              }else {


              throw new BizException(BizCodeEnum.COUPON_RECORD_LOCK_FAIL);
              }


              }


              • 消费者逻辑

                • 消息消费者端会根据消息的类型来匹配找一个RabbitHandler使用,第一个参数是消息的类型,后面两个是固定的

                • 发散思维:这里如果消息消费不成功就会一直入队,这样的话会导致队列爆满,如何解决上述问题呢? 可以用redis来存储消息,每次消息入队的时候就判断redis中这个消息的id的value是否为3,如果为3那么就不要入队,然后将这条消费不成功的记录插入数据库中,后续可以通过人工排查解决。

                @Slf4j
                @Component
                @RabbitListener(queues = "${mqconfig.coupon_release_queue}")
                public class CouponMQListener {


                @Autowired
                private CouponRecordService couponRecordService;


                @Autowired
                private RedissonClient redissonClient;


                /**
                *
                * 重复消费-幂等性
                *
                * 消费失败,重新入队后最大重试次数:
                * 如果消费失败,不重新入队,可以记录日志,然后插到数据库人工排查
                *
                * 消费者这块还有啥问题,大家可以先想下,然后给出解决方案
                *
                * @param recordMessage
                * @param message
                * @param channel
                * @throws IOException
                */
                @RabbitHandler
                public void releaseCouponRecord(CouponRecordMessage recordMessage, Message message, Channel channel) throws IOException {


                log.info("监听到消息:releaseCouponRecord消息内容:{}", recordMessage);
                long msgTag = message.getMessageProperties().getDeliveryTag();


                //核心处理逻辑搬到了CouponRecordServiceImpl中
                boolean flag = couponRecordService.releaseCouponRecord(recordMessage);


                //防止同个解锁任务并发进入;如果是串行消费不用加锁;加锁有利也有弊,看项目业务逻辑而定
                //Lock lock = redissonClient.getLock("lock:coupon_record_release:"+recordMessage.getTaskId());
                //lock.lock();
                try {
                if (flag) {
                //确认消息消费成功
                channel.basicAck(msgTag, false);
                }else {
                log.error("释放优惠券失败 flag=false,{}",recordMessage);
                channel.basicReject(msgTag,true);
                }


                } catch (IOException e) {
                log.error("释放优惠券记录异常:{},msg:{}",e,recordMessage);
                channel.basicReject(msgTag,true);
                }
                // finally {
                // lock.unlock();
                // }


                }
                }
                • CouponRecordServiceImpl消费者端的核心业务逻辑

                  /**
                  * 解锁优惠券记录
                  * 1)查询task工作单是否存在
                  * 2) 查询订单状态
                  * @param recordMessage
                  * @return
                  */
                  @Override
                  @Transactional(rollbackFor = Exception.class,propagation = Propagation.REQUIRED)
                  public boolean releaseCouponRecord(CouponRecordMessage recordMessage) {


                  //查询下task是否存
                  CouponTaskDO taskDO = couponTaskMapper.selectOne(new QueryWrapper<CouponTaskDO>().eq("id",recordMessage.getTaskId()));


                  if(taskDO==null){
                  log.warn("工作单不存,消息:{}",recordMessage);
                  return true;
                  }


                  //lock状态才处理
                  if(taskDO.getLockState().equalsIgnoreCase(StockTaskStateEnum.LOCK.name())){
                  //查询订单状态
                  JsonData jsonData = orderFeignSerivce.queryProductOrderState(recordMessage.getOutTradeNo());
                  if(jsonData.getCode()==0){
                  //正常响应,判断订单状态
                  String state = jsonData.getData().toString();
                  if(ProductOrderStateEnum.NEW.name().equalsIgnoreCase(state)){
                  //状态是NEW新建状态,则返回给消息队,列重新投递
                  log.warn("订单状态是NEW,返回给消息队列,重新投递:{}",recordMessage);
                  return false;
                  }
                  //如果是已经支付
                  if(ProductOrderStateEnum.PAY.name().equalsIgnoreCase(state)){
                  //如果已经支付,修改task状态为finish
                  taskDO.setLockState(StockTaskStateEnum.FINISH.name());
                  couponTaskMapper.update(taskDO,new QueryWrapper<CouponTaskDO>().eq("id",recordMessage.getTaskId()));
                  log.info("订单已经支付,修改库存锁定工作单FINISH状态:{}",recordMessage);
                  return true;
                  }
                  }
                  //订单不存在,或者订单被取消,确认消息,修改task状态为CANCEL,恢复优惠券使用记录为NEW
                  log.warn("订单不存在,或者订单被取消,确认消息,修改task状态为CANCEL,恢复优惠券使用记录为NEW,message:{}",recordMessage);
                  taskDO.setLockState(StockTaskStateEnum.CANCEL.name());


                  couponTaskMapper.update(taskDO,new QueryWrapper<CouponTaskDO>().eq("id",recordMessage.getTaskId()));
                  //恢复优惠券记录是NEW状态
                  couponRecordMapper.updateState(taskDO.getCouponRecordId(),CouponStateEnum.NEW.name());
                  return true;
                  }else {
                  log.warn("工作单状态不是LOCK,state={},消息体={}",taskDO.getLockState(),recordMessage);
                  return true;
                  }
                  }


                  算法vip班级永久班新春优惠价,全年最低价,加入后即可获取往期所有算法视频的录播课和后续算法班的永久直播授课,错过这次将再等一年,早加入早受益。

                  算法直播永久授课+超多资料福利,错过将再等一年~

                  奔跑的小梁,公众号:梁霖编程工具库算法训练营春节超低优惠价通知,错过这次再等一年!!!



                  文章转载自梁霖编程工具库,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

                  评论