diff --git a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/api/TransportTaskApi.java b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/api/TransportTaskApi.java index 7c446a98..44ba34a1 100644 --- a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/api/TransportTaskApi.java +++ b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/api/TransportTaskApi.java @@ -51,6 +51,12 @@ public interface TransportTaskApi { @Operation(summary = "根据 taskCode 查询任务") CommonResult getTaskByCode(@RequestParam("taskCode") String taskCode); + @GetMapping(PREFIX + "/getUnfinishedTaskByBiz") + @Operation(summary = "根据业务标识查询未完成任务") + CommonResult getUnfinishedTaskByBiz(@RequestParam("ownerService") String ownerService, + @RequestParam("bizType") String bizType, + @RequestParam("bizId") String bizId); + @GetMapping(PREFIX + "/getRunningTaskByMaterialId") @Operation(summary = "根据 materialId 和 ownerService 查询运行中任务") CommonResult> getRunningTaskByMaterialId(@RequestParam("materialId") Long materialId, diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/api/TransportTaskApiImpl.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/api/TransportTaskApiImpl.java index 178eaa32..e669265e 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/api/TransportTaskApiImpl.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/api/TransportTaskApiImpl.java @@ -60,6 +60,11 @@ public class TransportTaskApiImpl implements TransportTaskApi { return success(transportTaskService.getTaskInfoByCode(taskCode)); } + @Override + public CommonResult getUnfinishedTaskByBiz(String ownerService, String bizType, String bizId) { + return success(transportTaskService.getUnfinishedTaskInfoByBiz(ownerService, bizType, bizId)); + } + @Override public CommonResult> getRunningTaskByMaterialId(Long materialId, String ownerService) { return success(transportTaskService.getRunningTaskInfoByMaterialId(materialId, ownerService)); diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/dal/mysql/transporttask/TransportTaskMapper.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/dal/mysql/transporttask/TransportTaskMapper.java index b276a568..c5f268f9 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/dal/mysql/transporttask/TransportTaskMapper.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/dal/mysql/transporttask/TransportTaskMapper.java @@ -80,8 +80,9 @@ public interface TransportTaskMapper extends BaseMapperX { /** * 按业务归属查询未完结任务(幂等校验用) */ - default TransportTaskDO selectUnfinishedByBiz(String bizType, String bizId) { + default TransportTaskDO selectUnfinishedByBiz(String ownerService, String bizType, String bizId) { return selectOne(new LambdaQueryWrapperX() + .eq(TransportTaskDO::getOwnerService, ownerService) .eq(TransportTaskDO::getBizType, bizType) .eq(TransportTaskDO::getBizId, bizId) .lt(TransportTaskDO::getTaskStatus, TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode())); diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java index ddb09fb7..057bce9e 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskService.java @@ -105,6 +105,16 @@ public interface TransportTaskService { */ TaskInfoDTO getTaskInfoByCode(String taskCode); + /** + * 根据业务标识查询未完成任务。 + * + * @param ownerService 业务归属服务 + * @param bizType 业务类型 + * @param bizId 业务标识 + * @return 未完成任务,不存在时返回空 + */ + TaskInfoDTO getUnfinishedTaskInfoByBiz(String ownerService, String bizType, String bizId); + /** * 根据物料ID和业务归属查询运行中任务(状态小于75) * diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java index d46659ec..24178698 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskServiceImpl.java @@ -32,6 +32,8 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.springframework.validation.annotation.Validated; +import org.redisson.api.RLock; +import org.redisson.api.RedissonClient; import java.util.List; @@ -49,6 +51,8 @@ import static cn.code.nl.module.task.enums.ErrorCodeConstants.*; @Validated public class TransportTaskServiceImpl implements TransportTaskService { + private static final String TASK_BIZ_CREATE_LOCK_KEY = "task:biz:create:"; + @Resource private TransportTaskMapper transportTaskMapper; @@ -67,6 +71,9 @@ public class TransportTaskServiceImpl implements TransportTaskService { @Resource private CodeGenApi codeGenApi; + @Resource + private RedissonClient redissonClient; + @Override public Long createTransportTask(TransportTaskSaveReqVO createReqVO) { TransportTaskDO transportTask = BeanUtils.toBean(createReqVO, TransportTaskDO.class); @@ -80,6 +87,31 @@ public class TransportTaskServiceImpl implements TransportTaskService { bizNo = "{{#_ret}}", success = TASK_INFO_CREATE_BY_RPC) public Long createTransportTaskByRpc(TransportTaskCreateReqDTO reqDTO) { + if (StrUtil.isBlank(reqDTO.getOwnerService()) || StrUtil.isBlank(reqDTO.getBizType()) + || StrUtil.isBlank(reqDTO.getBizId())) { + return insertTransportTaskByRpc(reqDTO); + } + String lockKey = TASK_BIZ_CREATE_LOCK_KEY + reqDTO.getOwnerService() + ":" + + reqDTO.getBizType() + ":" + reqDTO.getBizId(); + RLock lock = redissonClient.getLock(lockKey); + lock.lock(); + try { + if (transportTaskMapper.selectUnfinishedByBiz(reqDTO.getOwnerService(), + reqDTO.getBizType(), reqDTO.getBizId()) != null) { + throw exception(TRANSPORT_TASK_EXISTS_UNFINISHED); + } + return insertTransportTaskByRpc(reqDTO); + } finally { + if (lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + } + + /** + * 插入 RPC 搬运任务。 + */ + private Long insertTransportTaskByRpc(TransportTaskCreateReqDTO reqDTO) { TransportTaskDO task = BeanUtils.toBean(reqDTO, TransportTaskDO.class); CodeGenerateReqDTO codeGenerateReqDTO = new CodeGenerateReqDTO(); codeGenerateReqDTO.setRuleCode("TASK_CODE"); @@ -218,6 +250,12 @@ public class TransportTaskServiceImpl implements TransportTaskService { return TransportTaskConvert.INSTANCE.convert(task); } + @Override + public TaskInfoDTO getUnfinishedTaskInfoByBiz(String ownerService, String bizType, String bizId) { + TransportTaskDO task = transportTaskMapper.selectUnfinishedByBiz(ownerService, bizType, bizId); + return TransportTaskConvert.INSTANCE.convert(task); + } + @Override public List getRunningTaskInfoByMaterialId(Long materialId, String ownerService) { List tasks = transportTaskMapper.selectRunningListByMaterialId(materialId, ownerService);