mirror of
http://192.168.1.205:9980/cf_devdept2/cf_imes_server.git
synced 2026-08-12 21:02:08 +08:00
生产单异步导入:1、修复死信队列ttl和对应任务状态修改;2、临时数据表移除无用状态;3、临时任务表新增文件字段;
This commit is contained in:
+5
-2
@@ -100,7 +100,7 @@ public class ChenfengRabbitMQAutoConfiguration {
|
|||||||
public Queue orderImportDefaultExcelQueue() {
|
public Queue orderImportDefaultExcelQueue() {
|
||||||
Map<String, Object> args = new HashMap<>(1);
|
Map<String, Object> args = new HashMap<>(1);
|
||||||
// x-dead-letter-exchange 这里声明当前队列绑定的死信交换机
|
// 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);
|
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();
|
return QueueBuilder.durable(RabbitMqConstants.ORDER_IMPORT_DEFAULT_EXCEL_QUEUE).withArguments(args).build();
|
||||||
@@ -153,15 +153,17 @@ public class ChenfengRabbitMQAutoConfiguration {
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* 生产单导入死信交换机声明
|
* 生产单导入死信交换机声明
|
||||||
|
*
|
||||||
* @return
|
* @return
|
||||||
*/
|
*/
|
||||||
@Bean
|
@Bean
|
||||||
public DirectExchange orderImportDeadLetterExchange(){
|
public DirectExchange orderImportDeadLetterExchange() {
|
||||||
return new DirectExchange(ORDER_IMPORT_DEAD_LETTER_EXCHANGE);
|
return new DirectExchange(ORDER_IMPORT_DEAD_LETTER_EXCHANGE);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 生产单导入死信队列声明
|
* 生产单导入死信队列声明
|
||||||
|
*
|
||||||
* @return
|
* @return
|
||||||
*/
|
*/
|
||||||
@Bean
|
@Bean
|
||||||
@@ -171,6 +173,7 @@ public class ChenfengRabbitMQAutoConfiguration {
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* 生产单导入死信队列、交换机、绑定声明
|
* 生产单导入死信队列、交换机、绑定声明
|
||||||
|
*
|
||||||
* @return
|
* @return
|
||||||
*/
|
*/
|
||||||
@Bean
|
@Bean
|
||||||
|
|||||||
+1
-1
@@ -25,7 +25,7 @@ public enum OrderImportStatusEnum {
|
|||||||
/**
|
/**
|
||||||
* 导入失败
|
* 导入失败
|
||||||
*/
|
*/
|
||||||
IMPORT_FAIL(2);
|
IMPORT_FAIL(20);
|
||||||
|
|
||||||
private int status;
|
private int status;
|
||||||
|
|
||||||
|
|||||||
+6
-1
@@ -30,7 +30,12 @@ public enum OrderImportTaskStatusEnum {
|
|||||||
/**
|
/**
|
||||||
* 消费失败
|
* 消费失败
|
||||||
*/
|
*/
|
||||||
CONSUME_FAIL(30);
|
CONSUME_FAIL(30),
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 消费超时
|
||||||
|
*/
|
||||||
|
CONSUME_TIMEOUT(40);
|
||||||
|
|
||||||
private int status;
|
private int status;
|
||||||
|
|
||||||
|
|||||||
+9
-25
@@ -13,31 +13,15 @@ import lombok.Getter;
|
|||||||
@AllArgsConstructor
|
@AllArgsConstructor
|
||||||
public enum OrderImportTmpDataTypeEnum {
|
public enum OrderImportTmpDataTypeEnum {
|
||||||
ORDER(0),
|
ORDER(0),
|
||||||
ROOM(1),
|
BODY(1),
|
||||||
BODY(2),
|
GROUP(2),
|
||||||
GROUP(3),
|
GOODS(3),
|
||||||
GOODS(4),
|
PLATE(4),
|
||||||
PLATE(6),
|
PART(5),
|
||||||
PART(7),
|
PART_EDGING(6),
|
||||||
PART_EDGING(8),
|
PART_HARDWARE(7),
|
||||||
PART_HARDWARE(9),
|
PART_ASSEMBLY(8),
|
||||||
PART_ASSEMBLY(10),
|
PART_REMARK(9);
|
||||||
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);
|
|
||||||
|
|
||||||
private int type;
|
private int type;
|
||||||
|
|
||||||
|
|||||||
+5
@@ -54,4 +54,9 @@ public class OrderImportTaskDO extends BaseDO {
|
|||||||
* 消费状态
|
* 消费状态
|
||||||
*/
|
*/
|
||||||
private Integer status;
|
private Integer status;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 导入文件名称
|
||||||
|
*/
|
||||||
|
private String fileName;
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-3
@@ -76,7 +76,7 @@ public class DefaultExcelOrderImportHandler implements AbstractExcelOrderImportH
|
|||||||
@Transactional(rollbackFor = Exception.class)
|
@Transactional(rollbackFor = Exception.class)
|
||||||
public void prepareAndConvert(OrderImportAsyncReqVO importAsyncReqVO, MultipartFile file) {
|
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());
|
String importLockKey = String.format(RedisKeyConstants.ORDER_IMPORT_LOCK_KEY, getUserOrganId());
|
||||||
// 检查机构导入锁
|
// 检查机构导入锁
|
||||||
boolean lockResult = redisLockUtil.lock(importLockKey, String.valueOf(taskId), chenfengCacheProperties.getLockTimeout());
|
boolean lockResult = redisLockUtil.lock(importLockKey, String.valueOf(taskId), chenfengCacheProperties.getLockTimeout());
|
||||||
@@ -102,7 +102,7 @@ public class DefaultExcelOrderImportHandler implements AbstractExcelOrderImportH
|
|||||||
|
|
||||||
Message message = MessageBuilder
|
Message message = MessageBuilder
|
||||||
.withBody(taskId.toString().getBytes())
|
.withBody(taskId.toString().getBytes())
|
||||||
.setHeader("x-message-ttl", 60000) // 设置 TTL 为 60000 毫秒(60 秒)
|
.setExpiration("60000") // 设置消息级 TTL(60 秒)
|
||||||
.build();
|
.build();
|
||||||
//如果处于事务中,待事务提交成功了再执行后续操作,防止事务未提交完成导致消费端查不到数据
|
//如果处于事务中,待事务提交成功了再执行后续操作,防止事务未提交完成导致消费端查不到数据
|
||||||
if (TransactionSynchronizationManager.isActualTransactionActive()) {
|
if (TransactionSynchronizationManager.isActualTransactionActive()) {
|
||||||
@@ -125,13 +125,14 @@ public class DefaultExcelOrderImportHandler implements AbstractExcelOrderImportH
|
|||||||
*
|
*
|
||||||
* @return taskId
|
* @return taskId
|
||||||
*/
|
*/
|
||||||
private Long createImportTask() {
|
private Long createImportTask(String fileName) {
|
||||||
// 创建导入任务
|
// 创建导入任务
|
||||||
OrderImportTaskDO orderImportTaskDO = OrderImportTaskDO.builder()
|
OrderImportTaskDO orderImportTaskDO = OrderImportTaskDO.builder()
|
||||||
.id((Long) snowFlakeGenerator.nextId(null))
|
.id((Long) snowFlakeGenerator.nextId(null))
|
||||||
.type(OrderImportTypeEnum.EXCEL.getType())
|
.type(OrderImportTypeEnum.EXCEL.getType())
|
||||||
.importStatus(OrderImportStatusEnum.IMPORT_NOT_YET.getStatus())
|
.importStatus(OrderImportStatusEnum.IMPORT_NOT_YET.getStatus())
|
||||||
.status(OrderImportTaskStatusEnum.CREATE_NOT_CONSUME.getStatus())
|
.status(OrderImportTaskStatusEnum.CREATE_NOT_CONSUME.getStatus())
|
||||||
|
.fileName(fileName)
|
||||||
.build();
|
.build();
|
||||||
orderImportTaskMapper.insert(orderImportTaskDO);
|
orderImportTaskMapper.insert(orderImportTaskDO);
|
||||||
return orderImportTaskDO.getId();
|
return orderImportTaskDO.getId();
|
||||||
|
|||||||
+4
-3
@@ -66,7 +66,7 @@ public class DefaultXmlOrderImportHandler implements AbstractXmlOrderImportHandl
|
|||||||
@Override
|
@Override
|
||||||
public void prepareAndConvert(OrderImportAsyncReqVO importAsyncReqVO, MultipartFile file) {
|
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());
|
String importLockKey = String.format(RedisKeyConstants.ORDER_IMPORT_LOCK_KEY, getUserOrganId());
|
||||||
// 检查机构导入锁
|
// 检查机构导入锁
|
||||||
boolean lockResult = redisLockUtil.lock(importLockKey, String.valueOf(taskId), chenfengCacheProperties.getLockTimeout());
|
boolean lockResult = redisLockUtil.lock(importLockKey, String.valueOf(taskId), chenfengCacheProperties.getLockTimeout());
|
||||||
@@ -89,7 +89,7 @@ public class DefaultXmlOrderImportHandler implements AbstractXmlOrderImportHandl
|
|||||||
|
|
||||||
Message message = MessageBuilder
|
Message message = MessageBuilder
|
||||||
.withBody(taskId.toString().getBytes())
|
.withBody(taskId.toString().getBytes())
|
||||||
.setHeader("x-message-ttl", 60000) // 设置 TTL 为 60000 毫秒(60 秒)
|
.setExpiration("60000") // 设置消息级 TTL(60 秒)
|
||||||
.build();
|
.build();
|
||||||
//如果处于事务中,待事务提交成功了再执行后续操作,防止事务未提交完成导致消费端查不到数据
|
//如果处于事务中,待事务提交成功了再执行后续操作,防止事务未提交完成导致消费端查不到数据
|
||||||
if (TransactionSynchronizationManager.isActualTransactionActive()) {
|
if (TransactionSynchronizationManager.isActualTransactionActive()) {
|
||||||
@@ -111,13 +111,14 @@ public class DefaultXmlOrderImportHandler implements AbstractXmlOrderImportHandl
|
|||||||
*
|
*
|
||||||
* @return taskId
|
* @return taskId
|
||||||
*/
|
*/
|
||||||
private Long createImportTask() {
|
private Long createImportTask(String fileName) {
|
||||||
// 创建导入任务
|
// 创建导入任务
|
||||||
OrderImportTaskDO orderImportTaskDO = OrderImportTaskDO.builder()
|
OrderImportTaskDO orderImportTaskDO = OrderImportTaskDO.builder()
|
||||||
.id((Long) snowFlakeGenerator.nextId(null))
|
.id((Long) snowFlakeGenerator.nextId(null))
|
||||||
.type(OrderImportTypeEnum.EXCEL.getType())
|
.type(OrderImportTypeEnum.EXCEL.getType())
|
||||||
.importStatus(OrderImportStatusEnum.IMPORT_NOT_YET.getStatus())
|
.importStatus(OrderImportStatusEnum.IMPORT_NOT_YET.getStatus())
|
||||||
.status(OrderImportTaskStatusEnum.CREATE_NOT_CONSUME.getStatus())
|
.status(OrderImportTaskStatusEnum.CREATE_NOT_CONSUME.getStatus())
|
||||||
|
.fileName(fileName)
|
||||||
.build();
|
.build();
|
||||||
orderImportTaskMapper.insert(orderImportTaskDO);
|
orderImportTaskMapper.insert(orderImportTaskDO);
|
||||||
return orderImportTaskDO.getId();
|
return orderImportTaskDO.getId();
|
||||||
|
|||||||
+1
-7
@@ -148,13 +148,7 @@ public class DefaultExcelOrderImportConsumer {
|
|||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
// 确认消息进入死信队列
|
// 确认消息进入死信队列
|
||||||
channel.basicNack(deliveryTag, false, false);
|
channel.basicNack(deliveryTag, false, false);
|
||||||
String importFailNotify = ORDER_IMPORT_FAIL.getMsg();
|
log.error(ORDER_IMPORT_FAIL.getMsg(), e);
|
||||||
// 更新状态为消费失败、导入失败
|
|
||||||
orderImportTaskMapper.update(new LambdaUpdateWrapper<OrderImportTaskDO>().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);
|
|
||||||
throw e;
|
throw e;
|
||||||
} finally {
|
} finally {
|
||||||
// 手动确认消息接收
|
// 手动确认消息接收
|
||||||
|
|||||||
+1
-7
@@ -144,13 +144,7 @@ public class DefaultXmlOrderImportConsumer {
|
|||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
// 确认消息进入死信队列
|
// 确认消息进入死信队列
|
||||||
channel.basicNack(deliveryTag, false, false);
|
channel.basicNack(deliveryTag, false, false);
|
||||||
String importFailNotify = ORDER_IMPORT_FAIL.getMsg();
|
log.error( ORDER_IMPORT_FAIL.getMsg(), e);
|
||||||
// 更新状态为消费失败、导入失败
|
|
||||||
orderImportTaskMapper.update(new LambdaUpdateWrapper<OrderImportTaskDO>().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);
|
|
||||||
throw e;
|
throw e;
|
||||||
} finally {
|
} finally {
|
||||||
// 手动确认消息接收
|
// 手动确认消息接收
|
||||||
|
|||||||
+50
-1
@@ -1,9 +1,15 @@
|
|||||||
package com.cf.imes.module.plan.service.orderImport.consumer;
|
package com.cf.imes.module.plan.service.orderImport.consumer;
|
||||||
|
|
||||||
import cn.hutool.core.util.ObjectUtil;
|
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.organ.core.context.OrganContextHolder;
|
||||||
import com.cf.imes.framework.redis.constants.RedisKeyConstants;
|
import com.cf.imes.framework.redis.constants.RedisKeyConstants;
|
||||||
import com.cf.imes.framework.redis.util.RedisLockUtil;
|
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 com.rabbitmq.client.Channel;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import org.springframework.amqp.rabbit.annotation.RabbitListener;
|
import org.springframework.amqp.rabbit.annotation.RabbitListener;
|
||||||
@@ -13,8 +19,11 @@ import org.springframework.stereotype.Component;
|
|||||||
|
|
||||||
import javax.annotation.Resource;
|
import javax.annotation.Resource;
|
||||||
import java.io.IOException;
|
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.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
|
@Resource
|
||||||
private RedisLockUtil redisLockUtil;
|
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)
|
@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<Map<String, Object>> xDeath) throws IOException {
|
||||||
log.info("====================【生产单导入收到死信消息:{}】====================", message);
|
log.info("====================【生产单导入收到死信消息:{}】====================", message);
|
||||||
// 手动确认消息
|
// 手动确认消息
|
||||||
channel.basicAck(deliveryTag, false);
|
channel.basicAck(deliveryTag, false);
|
||||||
|
|
||||||
|
String deadLetterReason = getDeadLetterReason(xDeath);
|
||||||
|
long taskId = Long.parseLong(message);
|
||||||
|
|
||||||
|
if(TTL_REASON.equals(deadLetterReason)) {
|
||||||
|
// 更新状态为消费超时
|
||||||
|
orderImportTaskMapper.update(new LambdaUpdateWrapper<OrderImportTaskDO>().eq(OrderImportTaskDO::getId, taskId)
|
||||||
|
.set(OrderImportTaskDO::getStatus, OrderImportTaskStatusEnum.CONSUME_FAIL.getStatus())
|
||||||
|
);
|
||||||
|
} else if (ERROR_REASON.equals(deadLetterReason)) {
|
||||||
|
// 更新状态为消费成功、导入失败
|
||||||
|
orderImportTaskMapper.update(new LambdaUpdateWrapper<OrderImportTaskDO>().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();
|
Long organId = OrganContextHolder.getOrganId();
|
||||||
if (ObjectUtil.isNotNull(organId)) {
|
if (ObjectUtil.isNotNull(organId)) {
|
||||||
@@ -42,4 +77,18 @@ public class OrderImportDeadLetterConsumer {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 获取进入死信队列原因
|
||||||
|
*
|
||||||
|
* @param xDeath
|
||||||
|
* @return
|
||||||
|
*/
|
||||||
|
private String getDeadLetterReason(List<Map<String, Object>> xDeath) {
|
||||||
|
if (xDeath == null || xDeath.isEmpty()) {
|
||||||
|
return "unknown";
|
||||||
|
}
|
||||||
|
// 通常取第一个 x-death 记录的原因(最新的)
|
||||||
|
Map<String, Object> firstDeath = xDeath.get(0);
|
||||||
|
return (String) firstDeath.get("reason");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user