Просмотр исходного кода

添加三个队列实现分队列管理

hxl13994548489 1 год назад
Родитель
Сommit
d2143e96e8

+ 11 - 1
src/main/java/com/ydtech/constants/RabbitMqEnum.java

@@ -11,7 +11,17 @@ public enum RabbitMqEnum {
     ORDER_QUEUE("order_push_queue", "队列"),
     ORDER_EXCHANGE("order_push_exchange", "交换机"),
 
-    ORDER_ROUTINGKEY("order_push_routingKey", "关键字");
+    ORDER_ROUTINGKEY("order_push_routingKey", "关键字"),
+
+    ORDER_QUEUE_FLAL("order_push_fail_queue", "消息失败队列"),
+    ORDER_EXCHANGE_FLAL("order_push_fail_exchange", "消息失败交换机"),
+
+    ORDER_ROUTINGKEY_FLAL("order_push_fail_routingKey", "消息失败关键字"),
+
+    ORDER_QUEUE_SIGN("order_push_sign_queue", "承包订单队列"),
+    ORDER_EXCHANGE_SIGN("order_push_sign_exchange", "承包订单交换机"),
+
+    ORDER_ROUTINGKEY_SIGN("order_push_sign_routingKey", "承包订单关键字");
 
     private final String code;
     private final String desc;

+ 72 - 0
src/main/java/com/ydtech/modules/order/orderThread/OrderFailRunnableThread.java

@@ -0,0 +1,72 @@
+package com.ydtech.modules.order.orderThread;
+
+import com.alibaba.fastjson.JSON;
+import com.ydtech.constants.RabbitMqEnum;
+import com.ydtech.modules.order.entity.po.InsOrderPushMsg;
+import com.ydtech.modules.order.entity.vo.PushDataVo;
+import com.ydtech.modules.order.service.InsOrderPushMsgService;
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.core.MessageProperties;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
+
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.UUID;
+
+import static org.springframework.amqp.core.MessageProperties.CONTENT_TYPE_TEXT_PLAIN;
+
+/**
+ *  @version
+ *  @author: hxl
+ *  @Date: 2024/11/5 15:50
+ *  @Description:  处理失败订单的的线程池
+ */
+public class OrderFailRunnableThread implements Runnable {
+
+    private RabbitTemplate rabbitTemplate;
+
+    private ArrayList<String> message =new ArrayList<>();
+
+    private InsOrderPushMsgService insOrderPushMsgService;
+
+    public OrderFailRunnableThread(RabbitTemplate rabbitTemplate, ArrayList<String> message, InsOrderPushMsgService insOrderPushMsgService) {
+        this.rabbitTemplate = rabbitTemplate;
+        this.message = message;
+        this.insOrderPushMsgService = insOrderPushMsgService;
+    }
+
+    public OrderFailRunnableThread() {
+    }
+
+    public OrderFailRunnableThread(InsOrderPushMsgService insOrderPushMsgService, ArrayList<String> message) {
+        this.message = message;
+        this.insOrderPushMsgService = insOrderPushMsgService;
+    }
+
+    @Override
+    public void run() {
+        //发送数据
+        this.message.forEach(item ->{
+            //修改执行相应业务逻辑
+            MessageProperties messageProperties = new MessageProperties();
+            messageProperties.setMessageId(UUID.randomUUID().toString());
+            messageProperties.setContentType(CONTENT_TYPE_TEXT_PLAIN);
+            messageProperties.setContentEncoding("UTF-8");
+            Message message_cj = new Message(item.getBytes(StandardCharsets.UTF_8), messageProperties);
+            PushDataVo pushDataVo = JSON.parseObject(item, PushDataVo.class);
+            //推送
+            InsOrderPushMsg insOrderPushMsg = new InsOrderPushMsg(null, pushDataVo.getInsAreaCompany().getId(), "0", new Date(), pushDataVo.toString());
+            try{
+                //更新狀態
+                insOrderPushMsgService.save(insOrderPushMsg);
+                this.rabbitTemplate.convertAndSend(RabbitMqEnum.ORDER_EXCHANGE_FLAL.getCode(), RabbitMqEnum.ORDER_ROUTINGKEY_FLAL.getCode(), message_cj);
+            }catch(Exception e){
+                e.getStackTrace();
+            }
+
+        });
+
+    }
+}
+

+ 72 - 0
src/main/java/com/ydtech/modules/order/orderThread/OrderSignRunnableThread.java

@@ -0,0 +1,72 @@
+package com.ydtech.modules.order.orderThread;
+
+import com.alibaba.fastjson.JSON;
+import com.ydtech.constants.RabbitMqEnum;
+import com.ydtech.modules.order.entity.po.InsOrderPushMsg;
+import com.ydtech.modules.order.entity.vo.PushDataVo;
+import com.ydtech.modules.order.service.InsOrderPushMsgService;
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.core.MessageProperties;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
+
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.UUID;
+
+import static org.springframework.amqp.core.MessageProperties.CONTENT_TYPE_TEXT_PLAIN;
+
+/**
+ *  @version
+ *  @author: hxl
+ *  @Date: 2024/11/5 15:50
+ *  @Description:  处理失败订单的的线程池
+ */
+public class OrderSignRunnableThread implements Runnable {
+
+    private RabbitTemplate rabbitTemplate;
+
+    private ArrayList<String> message =new ArrayList<>();
+
+    private InsOrderPushMsgService insOrderPushMsgService;
+
+    public OrderSignRunnableThread(RabbitTemplate rabbitTemplate, ArrayList<String> message, InsOrderPushMsgService insOrderPushMsgService) {
+        this.rabbitTemplate = rabbitTemplate;
+        this.message = message;
+        this.insOrderPushMsgService = insOrderPushMsgService;
+    }
+
+    public OrderSignRunnableThread() {
+    }
+
+    public OrderSignRunnableThread(InsOrderPushMsgService insOrderPushMsgService, ArrayList<String> message) {
+        this.message = message;
+        this.insOrderPushMsgService = insOrderPushMsgService;
+    }
+
+    @Override
+    public void run() {
+        //发送数据
+        this.message.forEach(item ->{
+            //修改执行相应业务逻辑
+            MessageProperties messageProperties = new MessageProperties();
+            messageProperties.setMessageId(UUID.randomUUID().toString());
+            messageProperties.setContentType(CONTENT_TYPE_TEXT_PLAIN);
+            messageProperties.setContentEncoding("UTF-8");
+            Message message_cj = new Message(item.getBytes(StandardCharsets.UTF_8), messageProperties);
+            PushDataVo pushDataVo = JSON.parseObject(item, PushDataVo.class);
+            //推送
+            InsOrderPushMsg insOrderPushMsg = new InsOrderPushMsg(null, pushDataVo.getInsAreaCompany().getId(), "0", new Date(), pushDataVo.toString());
+            try{
+                //更新狀態
+                insOrderPushMsgService.save(insOrderPushMsg);
+                this.rabbitTemplate.convertAndSend(RabbitMqEnum.ORDER_EXCHANGE_SIGN.getCode(), RabbitMqEnum.ORDER_ROUTINGKEY_SIGN.getCode(), message_cj);
+            }catch(Exception e){
+                e.getStackTrace();
+            }
+
+        });
+
+    }
+}
+

+ 4 - 2
src/main/java/com/ydtech/modules/order/service/impl/InsOrderPushMsgServiceImpl.java

@@ -11,7 +11,9 @@ import com.ydtech.modules.order.entity.InsOrders;
 import com.ydtech.modules.order.entity.dto.InsTaskImagesDto;
 import com.ydtech.modules.order.entity.po.InsOrderPushMsg;
 import com.ydtech.modules.order.entity.vo.PushDataVo;
+import com.ydtech.modules.order.orderThread.OrderFailRunnableThread;
 import com.ydtech.modules.order.orderThread.OrderRunnableThread;
+import com.ydtech.modules.order.orderThread.OrderSignRunnableThread;
 import com.ydtech.modules.order.service.InsAreaCompanyService;
 import com.ydtech.modules.order.service.InsOrderPushMsgService;
 import com.ydtech.modules.order.service.InsOrdersService;
@@ -100,7 +102,7 @@ public class InsOrderPushMsgServiceImpl extends ServiceImpl<InsOrderPushMsgMappe
             });
         }
         if(pushDataList!=null && !pushDataList.isEmpty()){
-            Thread thread = new Thread(new OrderRunnableThread(rabbitTemplate, pushDataList,this));
+            Thread thread = new Thread(new OrderFailRunnableThread(rabbitTemplate, pushDataList,this));
             thread.start();
         }
      }
@@ -115,7 +117,7 @@ public class InsOrderPushMsgServiceImpl extends ServiceImpl<InsOrderPushMsgMappe
         try{
             ArrayList<String>  pushDataList =new ArrayList<String>();
             pushDataList.add(getPushDataJsonString(item));
-            Thread thread = new Thread(new OrderRunnableThread(rabbitTemplate, pushDataList,this));
+            Thread thread = new Thread(new OrderSignRunnableThread(rabbitTemplate, pushDataList,this));
             thread.start();
         }catch(Exception e){
             log.error(e.getMessage());