From 9d670df28f5c35f47455e223e53620d2afb75d24 Mon Sep 17 00:00:00 2001 From: gaoqr <13665037151@163.com> Date: Thu, 13 Mar 2025 11:52:59 +0800 Subject: [PATCH] =?UTF-8?q?=E7=94=9F=E4=BA=A7=E5=8D=95=E5=BC=82=E6=AD=A5?= =?UTF-8?q?=E5=AF=BC=E5=85=A5=EF=BC=9A1=E3=80=81=E4=BF=AE=E5=A4=8D?= =?UTF-8?q?=E6=AD=BB=E4=BF=A1=E9=98=9F=E5=88=97ttl=E5=92=8C=E5=AF=B9?= =?UTF-8?q?=E5=BA=94=E4=BB=BB=E5=8A=A1=E7=8A=B6=E6=80=81=E4=BF=AE=E6=94=B9?= =?UTF-8?q?=EF=BC=9B2=E3=80=81=E4=B8=B4=E6=97=B6=E6=95=B0=E6=8D=AE?= =?UTF-8?q?=E8=A1=A8=E7=A7=BB=E9=99=A4=E6=97=A0=E7=94=A8=E7=8A=B6=E6=80=81?= =?UTF-8?q?=EF=BC=9B3=E3=80=81=E4=B8=B4=E6=97=B6=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E8=A1=A8=E6=96=B0=E5=A2=9E=E6=96=87=E4=BB=B6=E5=AD=97=E6=AE=B5?= =?UTF-8?q?=EF=BC=9B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../ChenfengRabbitMQAutoConfiguration.java | 7 ++- .../executor/enums/OrderImportStatusEnum.java | 2 +- .../enums/OrderImportTaskStatusEnum.java | 7 ++- .../enums/OrderImportTmpDataTypeEnum.java | 34 ++++--------- .../orderImport/OrderImportTaskDO.java | 5 ++ .../excel/DefaultExcelOrderImportHandler.java | 7 +-- .../xml/DefaultXmlOrderImportHandler.java | 7 +-- .../DefaultExcelOrderImportConsumer.java | 8 +-- .../DefaultXmlOrderImportConsumer.java | 8 +-- .../OrderImportDeadLetterConsumer.java | 51 ++++++++++++++++++- 10 files changed, 86 insertions(+), 50 deletions(-) diff --git a/cf-framework/cf-spring-boot-starter-mq/src/main/java/com/cf/imes/framework/mq/rabbitmq/config/ChenfengRabbitMQAutoConfiguration.java b/cf-framework/cf-spring-boot-starter-mq/src/main/java/com/cf/imes/framework/mq/rabbitmq/config/ChenfengRabbitMQAutoConfiguration.java index 36a70a543..3874b1fc7 100644 --- a/cf-framework/cf-spring-boot-starter-mq/src/main/java/com/cf/imes/framework/mq/rabbitmq/config/ChenfengRabbitMQAutoConfiguration.java +++ b/cf-framework/cf-spring-boot-starter-mq/src/main/java/com/cf/imes/framework/mq/rabbitmq/config/ChenfengRabbitMQAutoConfiguration.java @@ -100,7 +100,7 @@ public class ChenfengRabbitMQAutoConfiguration { public Queue orderImportDefaultExcelQueue() { Map args = new HashMap<>(1); // x-dead-letter-exchange 这里声明当前队列绑定的死信交换机 - args.put("x-dead-letter-exchange", ORDER_IMPORT_DEAD_LETTER_EXCHANGE); + args.put("x-dead-letter-exchange", ORDER_IMPORT_DEAD_LETTER_EXCHANGE); // 设置死信路由键 args.put("x-dead-letter-routing-key", ORDER_IMPORT_DEAD_LETTER_ROUTING_KEY); return QueueBuilder.durable(RabbitMqConstants.ORDER_IMPORT_DEFAULT_EXCEL_QUEUE).withArguments(args).build(); @@ -153,15 +153,17 @@ public class ChenfengRabbitMQAutoConfiguration { /** * 生产单导入死信交换机声明 + * * @return */ @Bean - public DirectExchange orderImportDeadLetterExchange(){ + public DirectExchange orderImportDeadLetterExchange() { return new DirectExchange(ORDER_IMPORT_DEAD_LETTER_EXCHANGE); } /** * 生产单导入死信队列声明 + * * @return */ @Bean @@ -171,6 +173,7 @@ public class ChenfengRabbitMQAutoConfiguration { /** * 生产单导入死信队列、交换机、绑定声明 + * * @return */ @Bean diff --git a/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportStatusEnum.java b/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportStatusEnum.java index 0e3ce3639..339116ef0 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportStatusEnum.java +++ b/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportStatusEnum.java @@ -25,7 +25,7 @@ public enum OrderImportStatusEnum { /** * 导入失败 */ - IMPORT_FAIL(2); + IMPORT_FAIL(20); private int status; diff --git a/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportTaskStatusEnum.java b/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportTaskStatusEnum.java index 10b6c86d9..9cdf46240 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportTaskStatusEnum.java +++ b/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportTaskStatusEnum.java @@ -30,7 +30,12 @@ public enum OrderImportTaskStatusEnum { /** * 消费失败 */ - CONSUME_FAIL(30); + CONSUME_FAIL(30), + + /** + * 消费超时 + */ + CONSUME_TIMEOUT(40); private int status; diff --git a/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportTmpDataTypeEnum.java b/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportTmpDataTypeEnum.java index d9a2c9ee4..3e8f2c847 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportTmpDataTypeEnum.java +++ b/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/OrderImportTmpDataTypeEnum.java @@ -13,31 +13,15 @@ import lombok.Getter; @AllArgsConstructor public enum OrderImportTmpDataTypeEnum { ORDER(0), - ROOM(1), - BODY(2), - GROUP(3), - GOODS(4), - PLATE(6), - PART(7), - PART_EDGING(8), - PART_HARDWARE(9), - PART_ASSEMBLY(10), - BLOCK_OUTLINE_POINT(11), - BLOCK_CUTTINGOUTLINE(12), - BLOCK_CUTTINGOUTLINE_POINT(13), - BLOCK_RAWOUTLINE_POINT(14), - BLOCK_HOLES(15), - BLOCK_HOLE(16), - BLOCK_GROOVES(17), - BLOCK_GROOVE(18), - BLOCK_GROOVE_OUTLINE_POINT(19), - BLOCK_GROOVE_POINTS_POINT(20), - BLOCK_GROOVE_ISLET(21), - BLOCK_GROOVE_ISLET_POINTS_POINT(22), - BLOCK_GROOVE_CORRELATIONS(23), - BLOCK_GROOVE_OFFSETLINE(24), - BLOCK_REMKAR(25), - PART_REMARK(26); + BODY(1), + GROUP(2), + GOODS(3), + PLATE(4), + PART(5), + PART_EDGING(6), + PART_HARDWARE(7), + PART_ASSEMBLY(8), + PART_REMARK(9); private int type; diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/dataobject/orderImport/OrderImportTaskDO.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/dataobject/orderImport/OrderImportTaskDO.java index 28e857d6a..640ae11bd 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/dataobject/orderImport/OrderImportTaskDO.java +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/dataobject/orderImport/OrderImportTaskDO.java @@ -54,4 +54,9 @@ public class OrderImportTaskDO extends BaseDO { * 消费状态 */ private Integer status; + + /** + * 导入文件名称 + */ + private String fileName; } diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/orderImport/handler/excel/DefaultExcelOrderImportHandler.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/orderImport/handler/excel/DefaultExcelOrderImportHandler.java index 0cbb864a1..89c729f44 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/orderImport/handler/excel/DefaultExcelOrderImportHandler.java +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/orderImport/handler/excel/DefaultExcelOrderImportHandler.java @@ -76,7 +76,7 @@ public class DefaultExcelOrderImportHandler implements AbstractExcelOrderImportH @Transactional(rollbackFor = Exception.class) public void prepareAndConvert(OrderImportAsyncReqVO importAsyncReqVO, MultipartFile file) { // 创建任务 - Long taskId = createImportTask(); + Long taskId = createImportTask(file.getOriginalFilename()); String importLockKey = String.format(RedisKeyConstants.ORDER_IMPORT_LOCK_KEY, getUserOrganId()); // 检查机构导入锁 boolean lockResult = redisLockUtil.lock(importLockKey, String.valueOf(taskId), chenfengCacheProperties.getLockTimeout()); @@ -102,7 +102,7 @@ public class DefaultExcelOrderImportHandler implements AbstractExcelOrderImportH Message message = MessageBuilder .withBody(taskId.toString().getBytes()) - .setHeader("x-message-ttl", 60000) // 设置 TTL 为 60000 毫秒(60 秒) + .setExpiration("60000") // 设置消息级 TTL(60 秒) .build(); //如果处于事务中,待事务提交成功了再执行后续操作,防止事务未提交完成导致消费端查不到数据 if (TransactionSynchronizationManager.isActualTransactionActive()) { @@ -125,13 +125,14 @@ public class DefaultExcelOrderImportHandler implements AbstractExcelOrderImportH * * @return taskId */ - private Long createImportTask() { + private Long createImportTask(String fileName) { // 创建导入任务 OrderImportTaskDO orderImportTaskDO = OrderImportTaskDO.builder() .id((Long) snowFlakeGenerator.nextId(null)) .type(OrderImportTypeEnum.EXCEL.getType()) .importStatus(OrderImportStatusEnum.IMPORT_NOT_YET.getStatus()) .status(OrderImportTaskStatusEnum.CREATE_NOT_CONSUME.getStatus()) + .fileName(fileName) .build(); orderImportTaskMapper.insert(orderImportTaskDO); return orderImportTaskDO.getId(); diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/orderImport/handler/xml/DefaultXmlOrderImportHandler.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/orderImport/handler/xml/DefaultXmlOrderImportHandler.java index 30db596c9..d589ea5d9 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/orderImport/handler/xml/DefaultXmlOrderImportHandler.java +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/orderImport/handler/xml/DefaultXmlOrderImportHandler.java @@ -66,7 +66,7 @@ public class DefaultXmlOrderImportHandler implements AbstractXmlOrderImportHandl @Override public void prepareAndConvert(OrderImportAsyncReqVO importAsyncReqVO, MultipartFile file) { // 创建任务 - Long taskId = createImportTask(); + Long taskId = createImportTask(file.getOriginalFilename()); String importLockKey = String.format(RedisKeyConstants.ORDER_IMPORT_LOCK_KEY, getUserOrganId()); // 检查机构导入锁 boolean lockResult = redisLockUtil.lock(importLockKey, String.valueOf(taskId), chenfengCacheProperties.getLockTimeout()); @@ -89,7 +89,7 @@ public class DefaultXmlOrderImportHandler implements AbstractXmlOrderImportHandl Message message = MessageBuilder .withBody(taskId.toString().getBytes()) - .setHeader("x-message-ttl", 60000) // 设置 TTL 为 60000 毫秒(60 秒) + .setExpiration("60000") // 设置消息级 TTL(60 秒) .build(); //如果处于事务中,待事务提交成功了再执行后续操作,防止事务未提交完成导致消费端查不到数据 if (TransactionSynchronizationManager.isActualTransactionActive()) { @@ -111,13 +111,14 @@ public class DefaultXmlOrderImportHandler implements AbstractXmlOrderImportHandl * * @return taskId */ - private Long createImportTask() { + private Long createImportTask(String fileName) { // 创建导入任务 OrderImportTaskDO orderImportTaskDO = OrderImportTaskDO.builder() .id((Long) snowFlakeGenerator.nextId(null)) .type(OrderImportTypeEnum.EXCEL.getType()) .importStatus(OrderImportStatusEnum.IMPORT_NOT_YET.getStatus()) .status(OrderImportTaskStatusEnum.CREATE_NOT_CONSUME.getStatus()) + .fileName(fileName) .build(); orderImportTaskMapper.insert(orderImportTaskDO); return orderImportTaskDO.getId(); diff --git a/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/DefaultExcelOrderImportConsumer.java b/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/DefaultExcelOrderImportConsumer.java index 4367e0c08..9ab3f896c 100644 --- a/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/DefaultExcelOrderImportConsumer.java +++ b/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/DefaultExcelOrderImportConsumer.java @@ -148,13 +148,7 @@ public class DefaultExcelOrderImportConsumer { } catch (Exception e) { // 确认消息进入死信队列 channel.basicNack(deliveryTag, false, false); - String importFailNotify = ORDER_IMPORT_FAIL.getMsg(); - // 更新状态为消费失败、导入失败 - orderImportTaskMapper.update(new LambdaUpdateWrapper().eq(OrderImportTaskDO::getId, taskId) - .set(OrderImportTaskDO::getStatus, OrderImportTaskStatusEnum.CONSUME_FAIL.getStatus()) - .set(OrderImportTaskDO::getImportStatus, OrderImportStatusEnum.IMPORT_FAIL.getStatus()) - .set(OrderImportTaskDO::getResult, importFailNotify)); - log.error(importFailNotify, e); + log.error(ORDER_IMPORT_FAIL.getMsg(), e); throw e; } finally { // 手动确认消息接收 diff --git a/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/DefaultXmlOrderImportConsumer.java b/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/DefaultXmlOrderImportConsumer.java index d8ac6fab8..405cb2135 100644 --- a/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/DefaultXmlOrderImportConsumer.java +++ b/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/DefaultXmlOrderImportConsumer.java @@ -144,13 +144,7 @@ public class DefaultXmlOrderImportConsumer { } catch (Exception e) { // 确认消息进入死信队列 channel.basicNack(deliveryTag, false, false); - String importFailNotify = ORDER_IMPORT_FAIL.getMsg(); - // 更新状态为消费失败、导入失败 - orderImportTaskMapper.update(new LambdaUpdateWrapper().eq(OrderImportTaskDO::getId, taskId) - .set(OrderImportTaskDO::getStatus, OrderImportTaskStatusEnum.CONSUME_FAIL.getStatus()) - .set(OrderImportTaskDO::getImportStatus, OrderImportStatusEnum.IMPORT_FAIL.getStatus()) - .set(OrderImportTaskDO::getResult, importFailNotify)); - log.error(importFailNotify, e); + log.error( ORDER_IMPORT_FAIL.getMsg(), e); throw e; } finally { // 手动确认消息接收 diff --git a/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/OrderImportDeadLetterConsumer.java b/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/OrderImportDeadLetterConsumer.java index a32e7e840..ab4b4aa36 100644 --- a/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/OrderImportDeadLetterConsumer.java +++ b/cf-module-prod-plan/cf-module-prod-plan-biz/src/main/java/com/cf/imes/module/plan/service/orderImport/consumer/OrderImportDeadLetterConsumer.java @@ -1,9 +1,15 @@ package com.cf.imes.module.plan.service.orderImport.consumer; import cn.hutool.core.util.ObjectUtil; +import com.baomidou.dynamic.datasource.annotation.DS; +import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper; import com.cf.imes.framework.organ.core.context.OrganContextHolder; import com.cf.imes.framework.redis.constants.RedisKeyConstants; import com.cf.imes.framework.redis.util.RedisLockUtil; +import com.cf.imes.module.executor.enums.OrderImportStatusEnum; +import com.cf.imes.module.executor.enums.OrderImportTaskStatusEnum; +import com.cf.imes.module.plan.dal.dataobject.orderImport.OrderImportTaskDO; +import com.cf.imes.module.plan.dal.mysql.orderImport.OrderImportTaskMapper; import com.rabbitmq.client.Channel; import lombok.extern.slf4j.Slf4j; import org.springframework.amqp.rabbit.annotation.RabbitListener; @@ -13,8 +19,11 @@ import org.springframework.stereotype.Component; import javax.annotation.Resource; import java.io.IOException; +import java.util.List; +import java.util.Map; import static com.cf.imes.framework.mq.rabbitmq.constant.RabbitMqConstants.ORDER_IMPORT_DEAD_LETTER_QUEUE; +import static com.cf.imes.module.plan.enums.ErrorCodeConstants.ORDER_IMPORT_FAIL; /** * 生产单导入死信队列消费 @@ -29,12 +38,38 @@ public class OrderImportDeadLetterConsumer { @Resource private RedisLockUtil redisLockUtil; + @Resource + private OrderImportTaskMapper orderImportTaskMapper; + + private static final String TTL_REASON = "expired"; + private static final String ERROR_REASON = "rejected"; + @RabbitListener(queues = ORDER_IMPORT_DEAD_LETTER_QUEUE) - public void receiveA(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { + @DS("imes_prod") + public void receive(String message, + Channel channel, + @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag, + @Header(value = "x-death", required = false) List> xDeath) throws IOException { log.info("====================【生产单导入收到死信消息:{}】====================", message); // 手动确认消息 channel.basicAck(deliveryTag, false); + String deadLetterReason = getDeadLetterReason(xDeath); + long taskId = Long.parseLong(message); + + if(TTL_REASON.equals(deadLetterReason)) { + // 更新状态为消费超时 + orderImportTaskMapper.update(new LambdaUpdateWrapper().eq(OrderImportTaskDO::getId, taskId) + .set(OrderImportTaskDO::getStatus, OrderImportTaskStatusEnum.CONSUME_FAIL.getStatus()) + ); + } else if (ERROR_REASON.equals(deadLetterReason)) { + // 更新状态为消费成功、导入失败 + orderImportTaskMapper.update(new LambdaUpdateWrapper().eq(OrderImportTaskDO::getId, taskId) + .set(OrderImportTaskDO::getStatus, OrderImportTaskStatusEnum.CONSUME_SUCCESS.getStatus()) + .set(OrderImportTaskDO::getImportStatus, OrderImportStatusEnum.IMPORT_FAIL.getStatus()) + .set(OrderImportTaskDO::getResult, ORDER_IMPORT_FAIL.getMsg())); + } + // 导入锁解锁 Long organId = OrganContextHolder.getOrganId(); if (ObjectUtil.isNotNull(organId)) { @@ -42,4 +77,18 @@ public class OrderImportDeadLetterConsumer { } } + /** + * 获取进入死信队列原因 + * + * @param xDeath + * @return + */ + private String getDeadLetterReason(List> xDeath) { + if (xDeath == null || xDeath.isEmpty()) { + return "unknown"; + } + // 通常取第一个 x-death 记录的原因(最新的) + Map firstDeath = xDeath.get(0); + return (String) firstDeath.get("reason"); + } }