12 KiB
12 KiB
handler_code 分发到对应 Handler 案例代码
1. 核心链路
MQ/RPC 收到 Task 事件
-> 读取 handlerCode
-> handlerMap.get(handlerCode)
-> 根据 eventType 调用 onTaskFinished / onTaskCancelled
旧版:
handle_class -> Class.forName -> StockAreaCallTubeTask#updateTaskStatus
新版:
handler_code -> handlerMap -> StockAreaCallTubeTaskHandler#onTaskFinished
2. 消息 DTO
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<String, Object> payload;
}
3. 统一 Handler 接口
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。
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<String, Object> 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<String> 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<String, Object> payload, String key) {
return payload == null || payload.get(key) == null ? null : String.valueOf(payload.get(key));
}
}
5. 另一个 Handler 示例
@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。
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<TaskBizCallbackHandler> handlers;
private final Map<String, TaskBizCallbackHandler> handlerMap = new HashMap<>();
public TaskBizCallbackHandlerRegistry(List<TaskBizCallbackHandler> 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;
}
}
注册结果类似:
LMS_STOCK_AREA_CALL_TUBE -> StockAreaCallTubeTaskHandler
LMS_SEND_AIR_SHAFT_AGV -> SendAirShaftAgvTaskHandler
7. 分发器
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 订阅:
Topic = TASK_EVENT_TOPIC
Tag = LMS
消费者:
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 伪代码:
@RocketMQMessageListener(
topic = "TASK_EVENT_TOPIC",
consumerGroup = "lms-task-event-consumer-group",
selectorExpression = "LMS"
)
@Component
@RequiredArgsConstructor
public class LmsTaskEventRocketConsumer implements RocketMQListener<TaskEventMessage> {
private final TaskBizCallbackDispatcher dispatcher;
@Override
public void onMessage(TaskEventMessage message) {
dispatcher.dispatch(message);
}
}
9. 幂等建议
生产环境建议在 dispatch 前做消费幂等:
1. 根据 eventId 查询消费日志
2. 如果已经 SUCCESS,直接 return
3. 如果没有,则插入 PROCESSING
4. 执行业务 handler
5. 更新消费日志为 SUCCESS
6. 异常时回滚或标记 FAILED,等待 MQ 重试
伪代码:
@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 为例:
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
-> 执行业务完成逻辑
最终,新版不是:
handle_class -> Class.forName -> updateTaskStatus
而是:
handler_code -> handlerMap -> onTaskFinished / onTaskCancelled