feat: task-job主表crud基础
This commit is contained in:
405
doc/handler-code-dispatch-example.md
Normal file
405
doc/handler-code-dispatch-example.md
Normal file
@@ -0,0 +1,405 @@
|
||||
# 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
|
||||
```
|
||||
Reference in New Issue
Block a user