Przeglądaj źródła

同步来客吧商户激活和商城订单数据

mac 2 lat temu
rodzic
commit
bf3ad626b8

+ 190 - 0
kxs-product/kxs-product-api/src/main/java/com/kxs/product/api/amqp/rabbit/RabbitKxsOrderQueueMQ.java

@@ -0,0 +1,190 @@
+package com.kxs.product.api.amqp.rabbit;
+
+import com.alibaba.fastjson.JSON;
+import com.alibaba.fastjson.JSONObject;
+import com.kxs.common.mq.enums.MQSendTypeEnum;
+import com.kxs.common.mq.model.AbstractMQ;
+import lombok.AllArgsConstructor;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import lombok.NoArgsConstructor;
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.core.MessageBuilder;
+import org.springframework.amqp.core.MessageDeliveryMode;
+import org.springframework.amqp.core.MessageProperties;
+
+import java.math.BigDecimal;
+import java.time.LocalDateTime;
+import java.util.UUID;
+
+/**
+ * mq 订单普通队列 配置
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-25
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@NoArgsConstructor
+@AllArgsConstructor
+public class RabbitKxsOrderQueueMQ extends AbstractMQ {
+
+
+    /**
+     * 订单交易队列
+     */
+    public static final String QUEUE_NAME = "QUEUE_KXS_ORDER_DIVISION";
+
+    /**
+     * 死队列名称
+     */
+    public static final String DEAD_QUEUE_NAME = null;
+
+    /**
+     * 内置msg 消息体定义
+     **/
+    private MsgEntity msgEntity;
+
+    /**
+     *  定义Msg消息载体
+     **/
+    @Data
+    public static class MsgEntity {
+
+        /**
+         * 订单ID
+         */
+        private String id;
+        /**
+         * 订单状态(0待付款,1已付款,2已完成,3已发货,4已退款)
+         */
+        private Integer status;
+        /**
+         * 订单创建时间
+         */
+        private String createDate;
+        /**
+         * 发货的机具券码
+         */
+        private String snNos;
+        /**
+         * 备注
+         */
+        private String remark;
+        /**
+         * 买入计数
+         */
+        private Integer buyCount;
+        /**
+         * 支付状态
+         */
+        private Integer payStatus;
+        /**
+         * 产品 ID
+         */
+        private Integer productId;
+        /**
+         * 发送状态
+         */
+        private Integer sendStatus;
+        /**
+         * 提货方式(1邮寄到付,2上门自提)
+         */
+        private Integer deliveryType;
+        /**
+         * 退款状态
+         */
+        private Integer refundStatus;
+        /**
+         * 支付方式(1支付宝,3余额,4储蓄金)
+         */
+        private Integer payMode;
+        /**
+         * 发送日期
+         */
+        private LocalDateTime sendDate;
+        /**
+         * 支付时间
+         */
+        private LocalDateTime payDate;
+        /**
+         * 地址
+         */
+        private String address;
+        /**
+         * 所在省市区
+         */
+        private String areas;
+        /**
+         * 支付总金额
+         */
+        private BigDecimal totalPrice;
+        /**
+         * 移动电话
+         */
+        private String mobile;
+        /**
+         * 真实姓名
+         */
+        private String realName;
+        /**
+         * 订单号
+         */
+        private String orderNo;
+        /**
+         * 用户 ID
+         */
+        private Integer userId;
+        /**
+         * 父订单 ID
+         */
+        private Integer parentOrderId;
+
+
+    }
+
+    @Override
+    public String getQueueName() {
+
+        return QUEUE_NAME;
+    }
+
+    @Override
+    public String getDeadQueueName() {
+
+        return DEAD_QUEUE_NAME;
+    }
+
+    @Override
+    public MQSendTypeEnum getMqType() {
+
+        return MQSendTypeEnum.QUEUE;
+    }
+
+    @Override
+    public Message toMessage() {
+        String message = JSONObject.toJSONString(msgEntity);
+        // 构建消息体
+        return MessageBuilder.withBody(message.getBytes())
+                .setContentType(MessageProperties.CONTENT_TYPE_TEXT_PLAIN)
+                .setDeliveryMode(MessageDeliveryMode.PERSISTENT)
+                .setMessageId(UUID.randomUUID().toString())
+                .build();
+    }
+
+    /**
+     * 构造发送消息
+     */
+    public static RabbitKxsOrderQueueMQ build(MsgEntity message){
+
+        return new RabbitKxsOrderQueueMQ(message);
+    }
+
+    /**
+     * 解析MQ消息, 一般用于接收MQ消息时
+     */
+    public static MsgEntity parse(String msg){
+        return JSON.parseObject(msg, MsgEntity.class);
+    }
+
+}

+ 59 - 0
kxs-product/kxs-product-biz/src/main/java/com/kxs/product/biz/mq/RabbitKxsOrderQueueListener.java

@@ -0,0 +1,59 @@
+package com.kxs.product.biz.mq;
+
+import com.alibaba.fastjson.JSON;
+import com.kxs.product.api.amqp.rabbit.RabbitKxsOrderQueueMQ;
+import com.kxs.product.api.amqp.rabbit.RabbitShopTimeoutQueueMQ;
+import com.kxs.product.api.model.KxsShopOrder;
+import com.kxs.product.biz.constant.enums.KxsShopEnum;
+import com.kxs.product.biz.service.KxsShopOrderService;
+import com.rabbitmq.client.Channel;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.rabbit.annotation.RabbitHandler;
+import org.springframework.amqp.rabbit.annotation.RabbitListener;
+import org.springframework.beans.BeanUtils;
+import org.springframework.stereotype.Component;
+
+import java.io.IOException;
+
+/**
+ * rabbit 队列侦听器 商品下单过期队列
+ *
+ * @author Pota1ovO
+ * @date 2024-05-07
+ */
+@Component
+@Slf4j
+@RequiredArgsConstructor
+public class RabbitKxsOrderQueueListener {
+
+	/**
+	 * 监听 老平台订单队列
+	 *
+	 * @param message 消息
+	 */
+	@RabbitListener(queues = RabbitKxsOrderQueueMQ.QUEUE_NAME, ackMode = "MANUAL")
+	@RabbitHandler
+	public void onMessage(String msg, Message message, Channel channel){
+		log.info("收到客小爽订单消息: " + RabbitKxsOrderQueueMQ.parse(msg));
+        try {
+			RabbitKxsOrderQueueMQ.MsgEntity parse = RabbitKxsOrderQueueMQ.parse(msg);
+			//处理百城千团订单,创建活动
+			if(parse.getProductId() == 100 && parse.getStatus() == 2){
+
+			}
+
+			channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
+        } catch (Exception e) {
+			log.error("客小爽订单消息消费失败:{}", JSON.toJSONString(msg), e);
+			try {
+				channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
+			} catch (IOException ex) {
+				log.error("客小爽订单消息手动确认失败" + ex.getMessage(), ex);
+			}
+        }
+    }
+
+
+}

+ 10 - 11
kxs-product/kxs-product-biz/src/main/java/com/kxs/product/biz/mq/RabbitShopQueueListener.java

@@ -1,5 +1,6 @@
 package com.kxs.product.biz.mq;
 
+import com.alibaba.fastjson.JSON;
 import com.kxs.product.api.amqp.rabbit.RabbitShopTimeoutQueueMQ;
 import com.kxs.product.api.model.KxsShopOrder;
 import com.kxs.product.biz.constant.enums.KxsShopEnum;
@@ -10,7 +11,6 @@ import lombok.extern.slf4j.Slf4j;
 import org.springframework.amqp.core.Message;
 import org.springframework.amqp.rabbit.annotation.RabbitHandler;
 import org.springframework.amqp.rabbit.annotation.RabbitListener;
-import org.springframework.beans.BeanUtils;
 import org.springframework.stereotype.Component;
 
 import java.io.IOException;
@@ -35,26 +35,25 @@ public class RabbitShopQueueListener {
 	@RabbitListener(queues = RabbitShopTimeoutQueueMQ.QUEUE_NAME, ackMode = "MANUAL")
 	@RabbitHandler
 	public void onMessage(String msg, Message message, Channel channel){
-		log.info("shop消费端Payload: " + RabbitShopTimeoutQueueMQ.parse(msg));
+		log.info("商品下单过期队列: " + RabbitShopTimeoutQueueMQ.parse(msg));
         try {
 			RabbitShopTimeoutQueueMQ.MsgEntity parse = RabbitShopTimeoutQueueMQ.parse(msg);
 			Integer id = parse.getId();
 			//根据订单号查询该订单是否付款成功,如果仍未付款成功,关闭订单
-			KxsShopOrder byId = kxsShopOrderService.getById(id);
-			if (byId != null && byId.getStatus().equals(KxsShopEnum.ORDER_NO_PAY.getType())){
-				KxsShopOrder kxsShopOrder = new KxsShopOrder();
-				BeanUtils.copyProperties(parse,kxsShopOrder);
-				//取消订单
+			KxsShopOrder order = kxsShopOrderService.getById(id);
+			if (order != null && order.getStatus().equals(KxsShopEnum.ORDER_NO_PAY.getType())){
+
 				//todo 这里只取消了我们平台的订单,如已发起支付渠道调用,仍需取消支付渠道
-				kxsShopOrderService.cancelOrder(kxsShopOrder);
+				order.setStatus(KxsShopEnum.ORDER_CANCEL.getType());
+				kxsShopOrderService.updateById(order);
 			}
 			channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
         } catch (Exception e) {
-			log.error(e.getMessage(), e);
+			log.error("商品下单过期队列消费失败:{}", JSON.toJSONString(msg), e);
 			try {
 				channel.basicReject(message.getMessageProperties().getDeliveryTag(), false);
 			} catch (IOException ex) {
-				log.error("进入死信队列失败" + ex.getMessage(), ex);
+				log.error("商品下单过期进入死信队列失败" + ex.getMessage(), ex);
 			}
         }
     }
@@ -67,7 +66,7 @@ public class RabbitShopQueueListener {
 	@RabbitListener(queues = RabbitShopTimeoutQueueMQ.DEAD_QUEUE_NAME)
 	@RabbitHandler
 	public void onDeadMessage(String msg, Message message, Channel channel){
-		log.info("死信队列Payload: " + RabbitShopTimeoutQueueMQ.parse(msg));
+		log.info("商品下单过期死信队列: " + RabbitShopTimeoutQueueMQ.parse(msg));
         try {
 			RabbitShopTimeoutQueueMQ.MsgEntity parse = RabbitShopTimeoutQueueMQ.parse(msg);
         } catch (Exception e) {

+ 2 - 2
kxs-stat/kxs-stat-biz/src/main/java/com/kxs/stat/biz/mq/RabbitLkbQueueListener.java

@@ -115,7 +115,7 @@ public class RabbitLkbQueueListener {
             try {
                 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
             } catch (IOException ex) {
-                log.error("来客吧商户激活消费失败" + ex.getMessage(), ex);
+                log.error("来客吧商户激活手动确认失败" + ex.getMessage(), ex);
             }
         }
     }
@@ -221,7 +221,7 @@ public class RabbitLkbQueueListener {
             try {
                 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
             } catch (IOException ex) {
-                log.error("来客吧交易数据消费失败" + ex.getMessage(), ex);
+                log.error("来客吧交易数据手动确认失败" + ex.getMessage(), ex);
             }
         }
     }