# handler_code 分发到对应 Handler 案例代码 ## 1. 核心链路 ```text MQ/RPC 收到 Task 事件 -> 读取 handlerCode -> handlerMap.get(handlerCode) -> 根据 eventType 调用 onTaskFinished / onTaskCancelled ``` 旧版: ```text handle_class -> Class.forName -> StockAreaCallTubeTask#updateTaskStatus ``` 新版: ```text handler_code -> handlerMap -> StockAreaCallTubeTaskHandler#onTaskFinished ``` ## 2. 消息 DTO ```java package cn.code.nl.module.lms.taskcallback.dto; import lombok.Data; import java.util.Map; @Data public class TaskEventMessage { private String eventId; private String eventType; // TASK_FINISHED / TASK_CANCELLED / TASK_EXECUTING private Long taskId; private String taskNo; private String ownerService; // LMS / WMS private String bizType; private String bizId; private String handlerCode; private Map payload; } ``` ## 3. 统一 Handler 接口 ```java package cn.code.nl.module.lms.taskcallback.handler; import cn.code.nl.module.lms.taskcallback.dto.TaskEventMessage; public interface TaskBizCallbackHandler { String getHandlerCode(); void onTaskFinished(TaskEventMessage message); void onTaskCancelled(TaskEventMessage message); default void onTaskExecuting(TaskEventMessage message) { } } ``` ## 4. StockAreaCallTubeTaskHandler 示例 对应旧版 `StockAreaCallTubeTask#updateTaskStatus`。 ```java package cn.code.nl.module.lms.taskcallback.handler.impl; import cn.code.nl.module.lms.taskcallback.dto.TaskEventMessage; import cn.code.nl.module.lms.taskcallback.handler.TaskBizCallbackHandler; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.springframework.transaction.annotation.Transactional; import java.util.Map; @Slf4j @Component @RequiredArgsConstructor public class StockAreaCallTubeTaskHandler implements TaskBizCallbackHandler { public static final String HANDLER_CODE = "LMS_STOCK_AREA_CALL_TUBE"; // private final BstIvtStockingivtService stockingivtService; // private final PdmBiSlittingproductionplanService slittingproductionplanService; // private final LmsTaskBizRecordService lmsTaskBizRecordService; @Override public String getHandlerCode() { return HANDLER_CODE; } @Override @Transactional(rollbackFor = Exception.class) public void onTaskFinished(TaskEventMessage message) { Map payload = message.getPayload(); String startPointCode = getString(payload, "startPointCode"); String endPointCode = getString(payload, "endPointCode"); String vehicleCode = getString(payload, "vehicleCode"); String toMaterial = getString(payload, "to_material"); log.info("备货区送纸管任务完成回调,taskId={}, bizId={}, start={}, end={}, vehicle={}", message.getTaskId(), message.getBizId(), startPointCode, endPointCode, vehicleCode); /* * 迁移旧版 FINISHED 分支逻辑: * 1. 查询起点备货位 * 2. 查询终点备货位 * 3. 终点设置为有货并写入 vehicle_code * 4. 起点清空 vehicle_code 并设置为空 * 5. 根据 to_material 恢复分切计划 is_paper_ok * 6. 更新 LMS 业务记录为完成 */ // BstIvtStockingivt startPoint = stockingivtService.getPointByCode(startPointCode, false); // BstIvtStockingivt endPoint = stockingivtService.getPointByCode(endPointCode, false); // endPoint.setIvtStatus("1"); // endPoint.setVehicleCode(vehicleCode); // stockingivtService.updateById(endPoint); // startPoint.setVehicleCode(""); // startPoint.setIvtStatus("0"); // stockingivtService.updateById(startPoint); // if (StrUtil.isNotBlank(toMaterial)) { // List tubes = Arrays.asList(toMaterial.split(",")); // slittingproductionplanService.restorePaperOkByMaterials(tubes); // } // lmsTaskBizRecordService.finish(message.getBizId()); } @Override @Transactional(rollbackFor = Exception.class) public void onTaskCancelled(TaskEventMessage message) { log.info("备货区送纸管任务取消回调,taskId={}, bizId={}", message.getTaskId(), message.getBizId()); /* * 迁移旧版 cancel 或取消分支逻辑: * 1. 恢复 LMS 业务记录状态 * 2. 释放目标点位预占 * 3. 恢复起点资源状态 * 4. 记录取消原因 */ // lmsTaskBizRecordService.cancel(message.getBizId(), "Task 服务通知取消"); } private String getString(Map payload, String key) { return payload == null || payload.get(key) == null ? null : String.valueOf(payload.get(key)); } } ``` ## 5. 另一个 Handler 示例 ```java @Component @RequiredArgsConstructor public class SendAirShaftAgvTaskHandler implements TaskBizCallbackHandler { public static final String HANDLER_CODE = "LMS_SEND_AIR_SHAFT_AGV"; @Override public String getHandlerCode() { return HANDLER_CODE; } @Override @Transactional(rollbackFor = Exception.class) public void onTaskFinished(TaskEventMessage message) { // 处理送气涨轴 AGV 任务完成逻辑 } @Override @Transactional(rollbackFor = Exception.class) public void onTaskCancelled(TaskEventMessage message) { // 处理送气涨轴 AGV 任务取消逻辑 } } ``` ## 6. Handler 注册表 Spring 启动时自动收集所有 `TaskBizCallbackHandler`。 ```java package cn.code.nl.module.lms.taskcallback.handler; import jakarta.annotation.PostConstruct; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import java.util.HashMap; import java.util.List; import java.util.Map; @Slf4j @Component public class TaskBizCallbackHandlerRegistry { private final List handlers; private final Map handlerMap = new HashMap<>(); public TaskBizCallbackHandlerRegistry(List handlers) { this.handlers = handlers; } @PostConstruct public void init() { for (TaskBizCallbackHandler handler : handlers) { String handlerCode = handler.getHandlerCode(); if (handlerMap.containsKey(handlerCode)) { throw new IllegalStateException("重复的任务业务处理器编码:" + handlerCode); } handlerMap.put(handlerCode, handler); log.info("注册任务业务处理器:{} -> {}", handlerCode, handler.getClass().getName()); } } public TaskBizCallbackHandler getRequiredHandler(String handlerCode) { TaskBizCallbackHandler handler = handlerMap.get(handlerCode); if (handler == null) { throw new IllegalArgumentException("未找到任务业务处理器,handlerCode=" + handlerCode); } return handler; } } ``` 注册结果类似: ```text LMS_STOCK_AREA_CALL_TUBE -> StockAreaCallTubeTaskHandler LMS_SEND_AIR_SHAFT_AGV -> SendAirShaftAgvTaskHandler ``` ## 7. 分发器 ```java package cn.code.nl.module.lms.taskcallback.service; import cn.code.nl.module.lms.taskcallback.dto.TaskEventMessage; import cn.code.nl.module.lms.taskcallback.handler.TaskBizCallbackHandler; import cn.code.nl.module.lms.taskcallback.handler.TaskBizCallbackHandlerRegistry; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @Slf4j @Service @RequiredArgsConstructor public class TaskBizCallbackDispatcher { private final TaskBizCallbackHandlerRegistry handlerRegistry; public void dispatch(TaskEventMessage message) { TaskBizCallbackHandler handler = handlerRegistry.getRequiredHandler(message.getHandlerCode()); switch (message.getEventType()) { case "TASK_FINISHED": handler.onTaskFinished(message); break; case "TASK_CANCELLED": handler.onTaskCancelled(message); break; case "TASK_EXECUTING": handler.onTaskExecuting(message); break; default: throw new IllegalArgumentException("不支持的任务事件类型:" + message.getEventType()); } } } ``` ## 8. MQ 消费者使用分发器 LMS 订阅: ```text Topic = TASK_EVENT_TOPIC Tag = LMS ``` 消费者: ```java package cn.code.nl.module.lms.taskcallback.mq; import cn.code.nl.module.lms.taskcallback.dto.TaskEventMessage; import cn.code.nl.module.lms.taskcallback.service.TaskBizCallbackDispatcher; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; @Slf4j @Component @RequiredArgsConstructor public class LmsTaskEventConsumer { private final TaskBizCallbackDispatcher dispatcher; public void onMessage(TaskEventMessage message) { if (!"LMS".equals(message.getOwnerService())) { log.warn("LMS 收到非 LMS 任务事件,ownerService={}, taskId={}", message.getOwnerService(), message.getTaskId()); return; } dispatcher.dispatch(message); } } ``` RocketMQ 伪代码: ```java @RocketMQMessageListener( topic = "TASK_EVENT_TOPIC", consumerGroup = "lms-task-event-consumer-group", selectorExpression = "LMS" ) @Component @RequiredArgsConstructor public class LmsTaskEventRocketConsumer implements RocketMQListener { private final TaskBizCallbackDispatcher dispatcher; @Override public void onMessage(TaskEventMessage message) { dispatcher.dispatch(message); } } ``` ## 9. 幂等建议 生产环境建议在 `dispatch` 前做消费幂等: ```text 1. 根据 eventId 查询消费日志 2. 如果已经 SUCCESS,直接 return 3. 如果没有,则插入 PROCESSING 4. 执行业务 handler 5. 更新消费日志为 SUCCESS 6. 异常时回滚或标记 FAILED,等待 MQ 重试 ``` 伪代码: ```java @Transactional(rollbackFor = Exception.class) public void dispatch(TaskEventMessage message) { // if (consumeLogService.isSuccess(message.getEventId())) { // return; // } // consumeLogService.markProcessing(message); TaskBizCallbackHandler handler = handlerRegistry.getRequiredHandler(message.getHandlerCode()); if ("TASK_FINISHED".equals(message.getEventType())) { handler.onTaskFinished(message); } else if ("TASK_CANCELLED".equals(message.getEventType())) { handler.onTaskCancelled(message); } else if ("TASK_EXECUTING".equals(message.getEventType())) { handler.onTaskExecuting(message); } else { throw new IllegalArgumentException("不支持的任务事件类型:" + message.getEventType()); } // consumeLogService.markSuccess(message.getEventId()); } ``` ## 10. 完整执行链路 以 `StockAreaCallTubeTask` 为例: ```text ACS 反馈任务完成 -> Task 服务接收反馈 -> Task 查询 task_task -> owner_service = LMS -> handler_code = LMS_STOCK_AREA_CALL_TUBE -> Task 发布 MQ:TASK_EVENT_TOPIC, Tag = LMS -> LMS 消费消息 -> LmsTaskEventConsumer.onMessage -> TaskBizCallbackDispatcher.dispatch -> TaskBizCallbackHandlerRegistry.getRequiredHandler("LMS_STOCK_AREA_CALL_TUBE") -> 找到 StockAreaCallTubeTaskHandler -> 调用 StockAreaCallTubeTaskHandler.onTaskFinished -> 执行业务完成逻辑 ``` 最终,新版不是: ```text handle_class -> Class.forName -> updateTaskStatus ``` 而是: ```text handler_code -> handlerMap -> onTaskFinished / onTaskCancelled ```