406 lines
12 KiB
Markdown
406 lines
12 KiB
Markdown
# 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<String, Object> 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<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 示例
|
||
|
||
```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<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;
|
||
}
|
||
}
|
||
```
|
||
|
||
注册结果类似:
|
||
|
||
```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<TaskEventMessage> {
|
||
|
||
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
|
||
```
|