fix: 任务并发控制
This commit is contained in:
@@ -1,15 +1,23 @@
|
||||
package cn.code.nl.module.lms.mq.consumer;
|
||||
|
||||
import cn.code.nl.framework.common.pojo.CommonResult;
|
||||
import cn.code.nl.framework.execute.core.AbstractTask;
|
||||
import cn.code.nl.framework.execute.core.TaskFactory;
|
||||
import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO;
|
||||
import cn.code.nl.module.task.api.TransportTaskApi;
|
||||
import cn.code.nl.module.task.dto.TaskInfoDTO;
|
||||
import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
|
||||
import cn.code.nl.module.task.message.TaskEventMessage;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
|
||||
import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
import org.redisson.api.RLock;
|
||||
import org.redisson.api.RedissonClient;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY;
|
||||
|
||||
/**
|
||||
* 监听任务状态变更,根据 handleCode 定位任务子类并执行完成/取消逻辑
|
||||
*
|
||||
@@ -27,12 +35,68 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
@Resource
|
||||
private TaskFactory taskFactory;
|
||||
|
||||
@Resource
|
||||
private TransportTaskApi transportTaskApi;
|
||||
|
||||
@Resource
|
||||
private RedissonClient redissonClient;
|
||||
|
||||
@Override
|
||||
public void onMessage(TaskEventMessage message) {
|
||||
String handleCode = message.getHandleCode();
|
||||
String eventType = message.getEventType();
|
||||
log.info("收到任务状态变更消息, handleCode={}, eventType={}, taskId={}", handleCode, eventType, message.getTaskId());
|
||||
|
||||
RLock lock = redissonClient.getLock(TASK_STATUS_CHANGE_LOCK_KEY + message.getTaskId());
|
||||
if (!lock.tryLock()) {
|
||||
log.warn("任务状态变更消息正在消费中,等待 MQ 重试, taskId={}, eventType={}", message.getTaskId(), eventType);
|
||||
throw new IllegalStateException("任务状态变更消息正在消费中");
|
||||
}
|
||||
try {
|
||||
if (!isTaskCallbackPending(message)) {
|
||||
return;
|
||||
}
|
||||
executeTaskHandler(message, handleCode, eventType);
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 判断任务是否处于当前事件对应的待业务处理状态
|
||||
*/
|
||||
private boolean isTaskCallbackPending(TaskEventMessage message) {
|
||||
String eventType = message.getEventType();
|
||||
String expectedStatus;
|
||||
if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) {
|
||||
expectedStatus = TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode();
|
||||
} else if (TaskEventMessage.EVENT_TYPE_CANCELLED.equals(eventType)) {
|
||||
expectedStatus = TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode();
|
||||
} else {
|
||||
log.warn("未知任务事件类型, eventType={}, taskId={}", eventType, message.getTaskId());
|
||||
return false;
|
||||
}
|
||||
|
||||
CommonResult<TaskInfoDTO> result = transportTaskApi.getTaskById(message.getTaskId());
|
||||
TaskInfoDTO taskInfo = result.getCheckedData();
|
||||
if (taskInfo == null) {
|
||||
log.warn("任务不存在,跳过任务状态变更消息, taskId={}, eventType={}", message.getTaskId(), eventType);
|
||||
return false;
|
||||
}
|
||||
if (!expectedStatus.equals(taskInfo.getTaskStatus())) {
|
||||
log.info("任务状态已处理,跳过重复消息, taskId={}, eventType={}, currentStatus={}, expectedStatus={}",
|
||||
message.getTaskId(), eventType, taskInfo.getTaskStatus(), expectedStatus);
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* 执行任务完成或取消业务处理器
|
||||
*/
|
||||
private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) {
|
||||
AbstractTask task = taskFactory.getTask(handleCode);
|
||||
if (task == null) {
|
||||
log.warn("未找到对应任务处理器, handleCode={}", handleCode);
|
||||
|
||||
@@ -3,9 +3,9 @@
|
||||
rocketmq:
|
||||
name-server: 192.168.81.193:9876
|
||||
producer:
|
||||
group: lms-producer-dev-group # 事务消息需要配置一样
|
||||
group: lms_producer_dev_group # 事务消息需要配置一样
|
||||
send-message-timeout: 3000
|
||||
consumer:
|
||||
lms-task-operate:
|
||||
group: lms-task-status-change-dev-group
|
||||
topic: lms-task-status-change-dev-topic
|
||||
group: lms_task_status_change_dev_group
|
||||
topic: lms_task_status_change_dev_topic
|
||||
@@ -34,4 +34,9 @@ public interface LogRecordConstants {
|
||||
String WMS_GROUP_PLATE = "组盘信息";
|
||||
String WMS_GROUP_PLATE_UPDATE = "修改组盘信息";
|
||||
String WMS_GROUP_PLATE_SUCCESS = "{{#loginUserNickname}} 更新了组盘信息: {{#group.vehicleCode}}";
|
||||
|
||||
// ======================= TASK 任务信息 =======================
|
||||
String TASK_INFO = "任务信息";
|
||||
String TASK_INFO_OPERATE_TYPE = "操作任务状态";
|
||||
String TASK_INFO_OPERATE_SUCCESS = "{{#loginUserNickname}}对任务操作了{{#operateName}}";
|
||||
}
|
||||
|
||||
@@ -14,6 +14,7 @@ import org.springframework.web.bind.annotation.RequestBody;
|
||||
import org.springframework.web.bind.annotation.RequestParam;
|
||||
|
||||
import jakarta.validation.Valid;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* Task 服务 RPC API 接口
|
||||
@@ -45,4 +46,14 @@ public interface TransportTaskApi {
|
||||
@GetMapping(PREFIX + "/getTaskByCode")
|
||||
@Operation(summary = "根据 taskCode 查询任务")
|
||||
CommonResult<TaskInfoDTO> getTaskByCode(@RequestParam("taskCode") String taskCode);
|
||||
|
||||
@GetMapping(PREFIX + "/getRunningTaskByMaterialId")
|
||||
@Operation(summary = "根据 materialId 和 ownerService 查询运行中任务")
|
||||
CommonResult<List<TaskInfoDTO>> getRunningTaskByMaterialId(@RequestParam("materialId") Long materialId,
|
||||
@RequestParam("ownerService") String ownerService);
|
||||
|
||||
@GetMapping(PREFIX + "/getRunningTaskByMaterialCode")
|
||||
@Operation(summary = "根据 materialCode 和 ownerService 查询运行中任务")
|
||||
CommonResult<List<TaskInfoDTO>> getRunningTaskByMaterialCode(@RequestParam("materialCode") String materialCode,
|
||||
@RequestParam("ownerService") String ownerService);
|
||||
}
|
||||
|
||||
@@ -91,6 +91,14 @@ public class TaskInfoDTO implements Serializable {
|
||||
* 载具编码2
|
||||
*/
|
||||
private String vehicleCode2;
|
||||
/**
|
||||
* 物料id
|
||||
*/
|
||||
private Long materialId;
|
||||
/**
|
||||
* 物料编码
|
||||
*/
|
||||
private String materialCode;
|
||||
/**
|
||||
* 车号
|
||||
*/
|
||||
|
||||
@@ -21,14 +21,6 @@ public class TransportTaskCreateReqDTO {
|
||||
@NotEmpty(message = "业务归属服务不能为空")
|
||||
private String ownerService;
|
||||
|
||||
@Schema(description = "业务类型", requiredMode = Schema.RequiredMode.REQUIRED)
|
||||
@NotEmpty(message = "业务类型不能为空")
|
||||
private String bizType;
|
||||
|
||||
@Schema(description = "业务侧标识", requiredMode = Schema.RequiredMode.REQUIRED)
|
||||
@NotEmpty(message = "业务侧标识不能为空")
|
||||
private String bizId;
|
||||
|
||||
@Schema(description = "业务回调处理器编码", requiredMode = Schema.RequiredMode.REQUIRED)
|
||||
@NotEmpty(message = "业务回调处理器编码不能为空")
|
||||
private String handleCode;
|
||||
@@ -59,6 +51,12 @@ public class TransportTaskCreateReqDTO {
|
||||
@Schema(description = "载具编码2")
|
||||
private String vehicleCode2;
|
||||
|
||||
@Schema(description = "物料id")
|
||||
private Long materialId;
|
||||
|
||||
@Schema(description = "物料编码")
|
||||
private String materialCode;
|
||||
|
||||
@Schema(description = "优先级")
|
||||
private String priority;
|
||||
|
||||
|
||||
@@ -0,0 +1,12 @@
|
||||
package cn.code.nl.module.task.message;
|
||||
|
||||
/**
|
||||
* 全局锁的key 或 前缀
|
||||
* @Author: liyongde
|
||||
* @Date: 2026/7/27 10:16
|
||||
*/
|
||||
public interface LockKeyConstants {
|
||||
|
||||
/** 任务状态变更消费锁前缀 */
|
||||
String TASK_STATUS_CHANGE_LOCK_KEY = "task:task-status-change:";
|
||||
}
|
||||
@@ -10,6 +10,8 @@ import jakarta.annotation.Resource;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import static cn.code.nl.framework.common.pojo.CommonResult.success;
|
||||
|
||||
/**
|
||||
@@ -51,4 +53,14 @@ public class TransportTaskApiImpl implements TransportTaskApi {
|
||||
public CommonResult<TaskInfoDTO> getTaskByCode(String taskCode) {
|
||||
return success(transportTaskService.getTaskInfoByCode(taskCode));
|
||||
}
|
||||
|
||||
@Override
|
||||
public CommonResult<List<TaskInfoDTO>> getRunningTaskByMaterialId(Long materialId, String ownerService) {
|
||||
return success(transportTaskService.getRunningTaskInfoByMaterialId(materialId, ownerService));
|
||||
}
|
||||
|
||||
@Override
|
||||
public CommonResult<List<TaskInfoDTO>> getRunningTaskByMaterialCode(String materialCode, String ownerService) {
|
||||
return success(transportTaskService.getRunningTaskInfoByMaterialCode(materialCode, ownerService));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,6 +5,8 @@ import cn.code.nl.module.task.dto.TaskInfoDTO;
|
||||
import org.mapstruct.Mapper;
|
||||
import org.mapstruct.factory.Mappers;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* 搬运任务 Convert
|
||||
*
|
||||
@@ -20,4 +22,9 @@ public interface TransportTaskConvert {
|
||||
* DO 转 RPC 全量信息 DTO(字段同名,零配置映射)
|
||||
*/
|
||||
TaskInfoDTO convert(TransportTaskDO bean);
|
||||
|
||||
/**
|
||||
* DO 列表转 RPC 全量信息 DTO 列表
|
||||
*/
|
||||
List<TaskInfoDTO> convertList(List<TransportTaskDO> list);
|
||||
}
|
||||
|
||||
@@ -40,11 +40,11 @@ public class TransportTaskDO extends BaseDO {
|
||||
*/
|
||||
private String ownerService;
|
||||
/**
|
||||
* 业务类型
|
||||
* 业务类型 todo: 暂时不用
|
||||
*/
|
||||
private String bizType;
|
||||
/**
|
||||
* 业务侧标识
|
||||
* 业务侧标识 todo: 暂时不用
|
||||
*/
|
||||
private String bizId;
|
||||
/**
|
||||
@@ -103,6 +103,14 @@ public class TransportTaskDO extends BaseDO {
|
||||
* 载具编码2
|
||||
*/
|
||||
private String vehicleCode2;
|
||||
/**
|
||||
* 物料id
|
||||
*/
|
||||
private Long materialId;
|
||||
/**
|
||||
* 物料编码
|
||||
*/
|
||||
private String materialCode;
|
||||
/**
|
||||
* 车号
|
||||
*/
|
||||
|
||||
@@ -87,6 +87,28 @@ public interface TransportTaskMapper extends BaseMapperX<TransportTaskDO> {
|
||||
.lt(TransportTaskDO::getTaskStatus, TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode()));
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据物料ID和业务归属查询运行中任务
|
||||
*/
|
||||
default List<TransportTaskDO> selectRunningListByMaterialId(Long materialId, String ownerService) {
|
||||
return selectList(new LambdaQueryWrapperX<TransportTaskDO>()
|
||||
.eq(TransportTaskDO::getMaterialId, materialId)
|
||||
.eq(TransportTaskDO::getOwnerService, ownerService)
|
||||
.lt(TransportTaskDO::getTaskStatus, TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode())
|
||||
.orderByDesc(TransportTaskDO::getTaskId));
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据物料编码和业务归属查询运行中任务
|
||||
*/
|
||||
default List<TransportTaskDO> selectRunningListByMaterialCode(String materialCode, String ownerService) {
|
||||
return selectList(new LambdaQueryWrapperX<TransportTaskDO>()
|
||||
.eq(TransportTaskDO::getMaterialCode, materialCode)
|
||||
.eq(TransportTaskDO::getOwnerService, ownerService)
|
||||
.lt(TransportTaskDO::getTaskStatus, TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode())
|
||||
.orderByDesc(TransportTaskDO::getTaskId));
|
||||
}
|
||||
|
||||
/**
|
||||
* 条件更新任务状态(CAS 抢占)
|
||||
*/
|
||||
|
||||
@@ -0,0 +1,16 @@
|
||||
package cn.code.nl.module.task.enums;
|
||||
|
||||
/**
|
||||
* 请求ACS 接口 定义常量
|
||||
* @Author: liyongde
|
||||
* @Date: 2026/7/27 9:45
|
||||
*/
|
||||
public interface AcsApiConstants {
|
||||
|
||||
/** 下发任务 */
|
||||
String ACS_TASK_API = "/acs-api/wms/issue-task";
|
||||
|
||||
/** 检测任务 */
|
||||
String ACS_OPERATE_CHECK_API = "/acs-api/wms/check-enable-operate";
|
||||
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
/**
|
||||
* 任务模块自身使用的枚举
|
||||
* @Author: liyongde
|
||||
* @Date: 2026/7/27 9:44
|
||||
*/
|
||||
package cn.code.nl.module.task.enums;
|
||||
@@ -24,6 +24,8 @@ import java.util.Set;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
import static cn.code.nl.module.task.enums.AcsApiConstants.ACS_TASK_API;
|
||||
|
||||
/**
|
||||
* 搬运任务下发管理器
|
||||
*/
|
||||
@@ -31,7 +33,6 @@ import java.util.stream.Collectors;
|
||||
@Component
|
||||
public class TransportTaskIssueManager {
|
||||
|
||||
private static final String ACS_TASK_API = "/acs-api/wms/task";
|
||||
private static final String ACS_SERVER_ADDRESS_CONFIG_SUFFIX = "-acs-server-address";
|
||||
|
||||
@Resource
|
||||
|
||||
@@ -14,6 +14,8 @@ import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import static cn.code.nl.module.task.enums.AcsApiConstants.ACS_OPERATE_CHECK_API;
|
||||
|
||||
/**
|
||||
* PC 端任务操作 ACS 校验管理器
|
||||
*/
|
||||
@@ -23,7 +25,6 @@ public class TransportTaskOperateCheckManager {
|
||||
|
||||
private static final String TASK_OPERATE_ENABLE_CONFIG_KEY = "task-operate-enable";
|
||||
private static final String TASK_OPERATE_ENABLE_VALUE = "1";
|
||||
private static final String ACS_OPERATE_CHECK_API = "/acs-api/wms/check-enable-operate";
|
||||
private static final String ACS_SERVER_ADDRESS_CONFIG_SUFFIX = "-acs-server-address";
|
||||
|
||||
@Resource
|
||||
|
||||
@@ -31,6 +31,8 @@ import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_AUT
|
||||
import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_STATUS_NOT_ALLOW;
|
||||
import static cn.code.nl.module.task.framework.common.util.TaskUtil.isAllowedFrom;
|
||||
|
||||
import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY;
|
||||
|
||||
/**
|
||||
* 搬运任务状态操作分发管理器
|
||||
*/
|
||||
@@ -70,7 +72,7 @@ public class TransportTaskOperationManager {
|
||||
* 按操作类型查路由表分发
|
||||
*/
|
||||
public void dispatchOperation(TransportTaskDO task, TaskOperationTypeEnum type, AcsFeedbackReqDTO reqDTO) {
|
||||
RLock lock = redissonClient.getLock(String.valueOf(task.getTaskId()));
|
||||
RLock lock = redissonClient.getLock(TASK_STATUS_CHANGE_LOCK_KEY + task.getTaskId());
|
||||
if (lock.tryLock()) {
|
||||
try {
|
||||
BiConsumer<TransportTaskDO, AcsFeedbackReqDTO> handler = operationHandlers.get(type);
|
||||
|
||||
@@ -31,10 +31,7 @@ public class TaskEventProducer {
|
||||
* @param message 事件消息
|
||||
*/
|
||||
public void publishEvent(TaskEventMessage message) {
|
||||
// 补齐 eventId
|
||||
if (message.getEventId() == null || message.getEventId().isEmpty()) {
|
||||
message.setEventId(UUID.randomUUID().toString());
|
||||
}
|
||||
message.setEventId(message.getTaskId().toString());
|
||||
|
||||
String destination = message.getOwnerService() + "_" + topic;
|
||||
try {
|
||||
|
||||
@@ -99,6 +99,24 @@ public interface TransportTaskService {
|
||||
*/
|
||||
TaskInfoDTO getTaskInfoByCode(String taskCode);
|
||||
|
||||
/**
|
||||
* 根据物料ID和业务归属查询运行中任务(状态小于75)
|
||||
*
|
||||
* @param materialId 物料ID
|
||||
* @param ownerService 业务归属服务
|
||||
* @return 运行中任务列表
|
||||
*/
|
||||
List<TaskInfoDTO> getRunningTaskInfoByMaterialId(Long materialId, String ownerService);
|
||||
|
||||
/**
|
||||
* 根据物料编码和业务归属查询运行中任务(状态小于75)
|
||||
*
|
||||
* @param materialCode 物料编码
|
||||
* @param ownerService 业务归属服务
|
||||
* @return 运行中任务列表
|
||||
*/
|
||||
List<TaskInfoDTO> getRunningTaskInfoByMaterialCode(String materialCode, String ownerService);
|
||||
|
||||
/**
|
||||
* 自动下发待下发任务到 ACS
|
||||
*/
|
||||
|
||||
@@ -2,6 +2,7 @@ package cn.code.nl.module.task.service.transporttask;
|
||||
|
||||
import cn.code.nl.framework.common.pojo.PageResult;
|
||||
import cn.code.nl.framework.common.util.object.BeanUtils;
|
||||
import cn.code.nl.framework.security.core.util.SecurityFrameworkUtils;
|
||||
import cn.code.nl.module.base.api.classstandard.ClassStandardApi;
|
||||
import cn.code.nl.module.task.controller.admin.transporttask.vo.TransportTaskOperateReqVO;
|
||||
import cn.code.nl.module.task.controller.admin.transporttask.vo.TransportTaskPageReqVO;
|
||||
@@ -21,6 +22,8 @@ import cn.code.nl.module.task.manage.TransportTaskOperateCheckManager;
|
||||
import cn.code.nl.module.task.manage.TransportTaskOperationManager;
|
||||
import cn.hutool.core.collection.CollUtil;
|
||||
import cn.hutool.core.util.StrUtil;
|
||||
import com.mzt.logapi.context.LogRecordContext;
|
||||
import com.mzt.logapi.starter.annotation.LogRecord;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Service;
|
||||
@@ -30,6 +33,7 @@ import org.springframework.validation.annotation.Validated;
|
||||
import java.util.List;
|
||||
|
||||
import static cn.code.nl.framework.common.exception.util.ServiceExceptionUtil.exception;
|
||||
import static cn.code.nl.module.system.enums.LogRecordConstants.*;
|
||||
import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_ALREADY_FINAL;
|
||||
import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_NOT_EXISTS;
|
||||
import static cn.code.nl.module.task.enums.ErrorCodeConstants.TRANSPORT_TASK_OPERATION_NOT_SUPPORTED;
|
||||
@@ -130,6 +134,10 @@ public class TransportTaskServiceImpl implements TransportTaskService {
|
||||
*/
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
@LogRecord(type = TASK_INFO,
|
||||
subType = TASK_INFO_OPERATE_TYPE,
|
||||
bizNo = "{{#reqVO.taskId}}",
|
||||
success = TASK_INFO_OPERATE_SUCCESS)
|
||||
public void operateTransportTask(TransportTaskOperateReqVO reqVO) {
|
||||
TransportTaskDO task = transportTaskMapper.selectById(reqVO.getTaskId());
|
||||
if (task == null) {
|
||||
@@ -155,6 +163,8 @@ public class TransportTaskServiceImpl implements TransportTaskService {
|
||||
reqDTO.setTaskId(task.getTaskId());
|
||||
reqDTO.setStatus(type.getCode());
|
||||
transportTaskOperationManager.dispatchOperation(task, type, reqDTO);
|
||||
LogRecordContext.putVariable("loginUserNickname", SecurityFrameworkUtils.getLoginUserNickname());
|
||||
LogRecordContext.putVariable("operateName", type.getName());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -186,6 +196,18 @@ public class TransportTaskServiceImpl implements TransportTaskService {
|
||||
return TransportTaskConvert.INSTANCE.convert(task);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<TaskInfoDTO> getRunningTaskInfoByMaterialId(Long materialId, String ownerService) {
|
||||
List<TransportTaskDO> tasks = transportTaskMapper.selectRunningListByMaterialId(materialId, ownerService);
|
||||
return TransportTaskConvert.INSTANCE.convertList(tasks);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<TaskInfoDTO> getRunningTaskInfoByMaterialCode(String materialCode, String ownerService) {
|
||||
List<TransportTaskDO> tasks = transportTaskMapper.selectRunningListByMaterialCode(materialCode, ownerService);
|
||||
return TransportTaskConvert.INSTANCE.convertList(tasks);
|
||||
}
|
||||
|
||||
/**
|
||||
* 查询待下发任务,并委托下发管理器按生产区域下发到 ACS
|
||||
*/
|
||||
|
||||
@@ -80,7 +80,7 @@ spring:
|
||||
rocketmq:
|
||||
name-server: 192.168.81.193:9876 # RocketMQ Namesrv
|
||||
producer:
|
||||
group: task-producer-dev-group # 生产者分组
|
||||
group: task_producer_dev_group # 生产者分组
|
||||
consumer:
|
||||
task:
|
||||
topic: task-status-change-dev-topic
|
||||
|
||||
@@ -1,15 +1,23 @@
|
||||
package cn.code.nl.module.wms.mq.consumer;
|
||||
|
||||
import cn.code.nl.framework.common.pojo.CommonResult;
|
||||
import cn.code.nl.framework.execute.core.AbstractTask;
|
||||
import cn.code.nl.framework.execute.core.TaskFactory;
|
||||
import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO;
|
||||
import cn.code.nl.module.task.api.TransportTaskApi;
|
||||
import cn.code.nl.module.task.dto.TaskInfoDTO;
|
||||
import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
|
||||
import cn.code.nl.module.task.message.TaskEventMessage;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
|
||||
import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
import org.redisson.api.RLock;
|
||||
import org.redisson.api.RedissonClient;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY;
|
||||
|
||||
/**
|
||||
* 监听任务状态变更,根据 handleCode 定位任务子类并执行完成/取消逻辑
|
||||
*
|
||||
@@ -27,12 +35,68 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
@Resource
|
||||
private TaskFactory taskFactory;
|
||||
|
||||
@Resource
|
||||
private TransportTaskApi transportTaskApi;
|
||||
|
||||
@Resource
|
||||
private RedissonClient redissonClient;
|
||||
|
||||
@Override
|
||||
public void onMessage(TaskEventMessage message) {
|
||||
String handleCode = message.getHandleCode();
|
||||
String eventType = message.getEventType();
|
||||
log.info("收到任务状态变更消息, handleCode={}, eventType={}, taskId={}", handleCode, eventType, message.getTaskId());
|
||||
|
||||
RLock lock = redissonClient.getLock(TASK_STATUS_CHANGE_LOCK_KEY + message.getTaskId());
|
||||
if (!lock.tryLock()) {
|
||||
log.warn("任务状态变更消息正在消费中,等待 MQ 重试, taskId={}, eventType={}", message.getTaskId(), eventType);
|
||||
throw new IllegalStateException("任务状态变更消息正在消费中......");
|
||||
}
|
||||
try {
|
||||
if (!isTaskCallbackPending(message)) {
|
||||
return;
|
||||
}
|
||||
executeTaskHandler(message, handleCode, eventType);
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 判断任务是否处于当前事件对应的待业务处理状态
|
||||
*/
|
||||
private boolean isTaskCallbackPending(TaskEventMessage message) {
|
||||
String eventType = message.getEventType();
|
||||
String expectedStatus;
|
||||
if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) {
|
||||
expectedStatus = TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode();
|
||||
} else if (TaskEventMessage.EVENT_TYPE_CANCELLED.equals(eventType)) {
|
||||
expectedStatus = TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode();
|
||||
} else {
|
||||
log.warn("未知任务事件类型, eventType={}, taskId={}", eventType, message.getTaskId());
|
||||
return false;
|
||||
}
|
||||
|
||||
CommonResult<TaskInfoDTO> result = transportTaskApi.getTaskById(message.getTaskId());
|
||||
TaskInfoDTO taskInfo = result.getCheckedData();
|
||||
if (taskInfo == null) {
|
||||
log.warn("任务不存在,跳过任务状态变更消息, taskId={}, eventType={}", message.getTaskId(), eventType);
|
||||
return false;
|
||||
}
|
||||
if (!expectedStatus.equals(taskInfo.getTaskStatus())) {
|
||||
log.info("任务状态已处理,跳过重复消息, taskId={}, eventType={}, currentStatus={}, expectedStatus={}",
|
||||
message.getTaskId(), eventType, taskInfo.getTaskStatus(), expectedStatus);
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* 执行任务完成或取消业务处理器
|
||||
*/
|
||||
private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) {
|
||||
AbstractTask task = taskFactory.getTask(handleCode);
|
||||
if (task == null) {
|
||||
log.warn("未找到对应任务处理器, handleCode={}", handleCode);
|
||||
|
||||
@@ -3,9 +3,9 @@
|
||||
rocketmq:
|
||||
name-server: 192.168.81.193:9876
|
||||
producer:
|
||||
group: wms-producer-dev-group # 事务消息需要配置一样
|
||||
group: wms_producer_dev_group # 事务消息需要配置一样
|
||||
send-message-timeout: 3000
|
||||
consumer:
|
||||
wms-task-operate:
|
||||
group: wms-task-status-change-dev-group
|
||||
topic: wms-task-status-change-dev-topic
|
||||
group: wms_task_status_change_dev_group
|
||||
topic: wms_task_status_change_dev_topic
|
||||
@@ -3,9 +3,14 @@
|
||||
rocketmq:
|
||||
name-server: 192.168.81.193:9876
|
||||
producer:
|
||||
group: nl-producer-dev-group
|
||||
group: nl_producer_dev_group
|
||||
send-message-timeout: 3000
|
||||
consumer:
|
||||
wms-task-operate:
|
||||
group: wms-task-status-change-dev-group
|
||||
topic: wms-task-status-change-dev-topic
|
||||
group: wms_task_status_change_dev_group
|
||||
topic: wms_task_status_change_dev_topic
|
||||
lms-task-operate:
|
||||
group: lms_task_status_change_dev_group
|
||||
topic: lms_task_status_change_dev_topic
|
||||
task:
|
||||
topic: task_status_change_dev_topic
|
||||
16
nl-server/src/main/resources/application-mq-test.yaml
Normal file
16
nl-server/src/main/resources/application-mq-test.yaml
Normal file
@@ -0,0 +1,16 @@
|
||||
--- #################### MQ 消息队列相关配置 ####################
|
||||
# rocketmq 配置项,对应 RocketMQProperties 配置类
|
||||
rocketmq:
|
||||
name-server: 192.168.81.193:9876
|
||||
producer:
|
||||
group: nl_producer_dev_group
|
||||
send-message-timeout: 3000
|
||||
consumer:
|
||||
wms-task-operate:
|
||||
group: wms_task_status_change_dev_group
|
||||
topic: wms_task_status_change_dev_topic
|
||||
lms-task-operate:
|
||||
group: lms_task_status_change_dev_group
|
||||
topic: lms_task_status_change_dev_topic
|
||||
task:
|
||||
topic: task_status_change_dev_topic
|
||||
@@ -49,20 +49,20 @@ spring:
|
||||
master:
|
||||
url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
|
||||
username: root
|
||||
password: root
|
||||
password: root123
|
||||
slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改
|
||||
lazy: true # 开启懒加载,保证启动速度
|
||||
url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
|
||||
username: root
|
||||
password: root
|
||||
password: root123
|
||||
|
||||
# Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优
|
||||
data:
|
||||
redis:
|
||||
host: 192.168.81.193 # 地址
|
||||
host: 127.0.0.1 # 地址
|
||||
port: 6379 # 端口
|
||||
database: 1 # 数据库索引
|
||||
password: redis123
|
||||
# password: redis123
|
||||
|
||||
--- #################### 定时任务相关配置 ####################
|
||||
|
||||
|
||||
@@ -14,5 +14,7 @@ export const ACTION_ICON = {
|
||||
BOOK: 'lucide:book',
|
||||
AUDIT: 'lucide:file-check',
|
||||
SEND: 'lucide:send',
|
||||
CIRCLE_CHECK: 'lucide:circle-check',
|
||||
CIRCLE_X: 'lucide:circle-x',
|
||||
CANCEL: 'lucide:ban',
|
||||
};
|
||||
|
||||
@@ -179,24 +179,24 @@ const [Grid, gridApi] = useVbenVxeGrid({
|
||||
label: '下发',
|
||||
type: 'link',
|
||||
icon: ACTION_ICON.SEND,
|
||||
auth: ['task:transport-task:operate'],
|
||||
auth: ['task:transport-task:issue'],
|
||||
disabled: row.taskStatus !== '40',
|
||||
onClick: handleOperate.bind(null, row, 'ISSUE', '下发任务'),
|
||||
},
|
||||
{
|
||||
label: '完成',
|
||||
type: 'link',
|
||||
icon: ACTION_ICON.SEND,
|
||||
icon: ACTION_ICON.CIRCLE_CHECK,
|
||||
auth: ['task:transport-task:finish'],
|
||||
disabled: ['79', '89', '99'].includes(row.taskStatus),
|
||||
disabled: ['79', '85', '89', '99'].includes(row.taskStatus),
|
||||
onClick: handleOperate.bind(null, row, 'FINISHED', '完成任务'),
|
||||
},
|
||||
{
|
||||
label: '取消',
|
||||
type: 'link',
|
||||
icon: ACTION_ICON.CANCEL,
|
||||
icon: ACTION_ICON.CIRCLE_X,
|
||||
auth: ['task:transport-task:cancel'],
|
||||
disabled: ['79', '89', '99'].includes(row.taskStatus),
|
||||
disabled: ['75', '79', '89', '99'].includes(row.taskStatus),
|
||||
onClick: handleOperate.bind(null, row, 'CANCELLED', '取消任务'),
|
||||
},
|
||||
{
|
||||
|
||||
Reference in New Issue
Block a user