Files
huachuang/doc/task-event-routing-to-business-handler-design.md

631 lines
15 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# 新版本任务完成/取消如何精准路由到具体业务 Handler
## 1. 问题背景
旧版兰州 LMS 中,任务类例如:
```text
StockAreaCallTubeTask
备货区送纸管到机械手旁边的备货区 - AGV任务
```
它继承 `AbstractAcsTask`,创建任务时写入:
```text
handle_class = org.nl.b_lms.sch.tasks.slitter.StockAreaCallTubeTask
```
当 ACS 请求任务完成时,旧系统会根据任务表里的 `handle_class` 找到对应 Spring Bean然后执行
```text
StockAreaCallTubeTask#updateTaskStatus(taskObj, status)
```
所以旧版的路由链路是:
```text
ACS 反馈完成
-> LMS 接收反馈
-> 查询 SCH_BASE_Task
-> 读取 handle_class
-> Class.forName(handle_class)
-> Spring 获取 StockAreaCallTubeTask Bean
-> 调用 updateTaskStatus
```
新版本中,任务模块独立为 `Task` 服务。`Task` 服务负责接收 ACS 反馈,但 `StockAreaCallTubeTask` 这类业务逻辑应该放在 `LMS` 服务中,而不是放在 `Task` 服务中。
因此新版本要解决两个路由问题:
1. Task 服务如何知道这个任务完成后应该通知 LMS 还是 WMS
2. LMS/WMS 收到消息后,如何精准定位到原来类似 `StockAreaCallTubeTask` 的业务逻辑?
## 2. 核心结论
新版本不再依赖 Java 类名 `handle_class` 跨服务反射,而是使用:
```text
owner_service + handler_code + event_type
```
来完成路由。
对应关系:
| 旧版 | 新版 |
| --- | --- |
| `handle_class = StockAreaCallTubeTask` | `handler_code = LMS_STOCK_AREA_CALL_TUBE` |
| LMS 内反射 Java 类 | LMS 内部 handler 注册表分发 |
| `updateTaskStatus(taskObj, status)` | `onTaskFinished / onTaskCancelled` |
| 单体内方法调用 | Task 服务发 MQ 或 RPC 到业务服务 |
新版本完整链路:
```text
ACS 反馈任务完成
-> Task 服务接收反馈
-> Task 查询任务记录
-> 得到 owner_service = LMS
-> 得到 handler_code = LMS_STOCK_AREA_CALL_TUBE
-> Task 发布 MQTaskFinishedEvent
-> LMS 消费属于自己的任务事件
-> LMS 根据 handler_code 找到 StockAreaCallTubeCallbackHandler
-> 执行 onTaskFinished
-> LMS 处理完成后回调 Task 或发送处理结果事件
-> Task 更新任务最终状态
```
## 3. 任务创建时必须保存路由信息
新版本能不能精准定位业务 handler关键取决于创建任务时保存的信息是否完整。
`StockAreaCallTubeTask` 为例,旧版创建任务保存:
```text
handle_class = org.nl.b_lms.sch.tasks.slitter.StockAreaCallTubeTask
```
新版本创建任务时LMS 应调用 Task 服务并传入:
```json
{
"ownerService": "LMS",
"bizType": "STOCK_AREA_CALL_TUBE",
"bizId": "LMS业务记录ID",
"handlerCode": "LMS_STOCK_AREA_CALL_TUBE",
"externalSystem": "ACS",
"externalTaskType": "3",
"startPointCode": "备货区起点",
"endPointCode": "机械手旁备货区终点",
"vehicleCode": "纸管载具号",
"requestPayload": {
"to_material": "xxx,xxx",
"product_area": "BLK"
}
}
```
其中最关键的是:
| 字段 | 作用 |
| --- | --- |
| `ownerService` | 告诉 Task 服务后续业务回调发给谁,例如 LMS/WMS |
| `handlerCode` | 告诉业务服务内部该由哪个 handler 处理 |
| `bizType` | 业务类型,用于查询、统计、权限和辅助路由 |
| `bizId` | 业务侧记录 IDhandler 可用它查询业务数据 |
| `requestPayload` | 保存旧任务中 `request_param` 类似的扩展数据 |
## 4. MQ Topic 应该如何设计
### 4.1 推荐方案:统一 Topic + owner_service 作为 Tag
推荐使用一个统一 Topic
```text
TASK_EVENT_TOPIC
```
不同业务服务通过 Tag 区分:
```text
Tag = LMS
Tag = WMS
```
消息体里包含完整路由字段:
```json
{
"eventId": "EVT-10001",
"eventType": "TASK_FINISHED",
"taskId": 10001,
"taskNo": "T202607130001",
"ownerService": "LMS",
"bizType": "STOCK_AREA_CALL_TUBE",
"bizId": "LMS-BIZ-10001",
"handlerCode": "LMS_STOCK_AREA_CALL_TUBE",
"externalSystem": "ACS",
"externalTaskNo": "ACS-001",
"payload": {
"to_material": "A,B,C"
}
}
```
LMS 只订阅:
```text
Topic = TASK_EVENT_TOPIC
Tag = LMS
```
WMS 只订阅:
```text
Topic = TASK_EVENT_TOPIC
Tag = WMS
```
优点:
- Topic 数量少;
- Task 事件模型统一;
- 新增业务服务只需要新增 Tag
- 消息结构统一,便于日志、监控、重试。
### 4.2 备选方案:按服务拆 Topic
也可以拆成:
```text
LMS_TASK_EVENT_TOPIC
WMS_TASK_EVENT_TOPIC
```
这种方式简单直接,但 Topic 会随着服务增多而增多。适合服务数量固定、消息隔离要求强的场景。
### 4.3 不推荐按任务类型拆 Topic
不推荐:
```text
STOCK_AREA_CALL_TUBE_TOPIC
WMS_INBOUND_TOPIC
...
```
因为任务类型会越来越多Topic 管理会失控。任务类型应该放在消息体的 `bizType``handlerCode` 中,而不是通过 Topic 区分。
## 5. Task 服务发布事件的逻辑
当 ACS 反馈任务完成时Task 服务处理步骤:
```text
1. 接收 ACS 反馈
2. 根据 taskId 查询 task_task
3. 幂等校验 external_event_id
4. 将任务状态改为 FINISHED_PENDING_CALLBACK
5. 根据 owner_service 决定消息 Tag
6. 发布 TASK_FINISHED 事件
7. 等待 LMS/WMS 业务处理结果
```
伪代码:
```java
public void handleAcsFeedback(AcsFeedbackRequest req) {
TaskDO task = taskMapper.selectById(req.getTaskId());
checkFeedbackIdempotent(req);
if (req.isFinished()) {
task.setStatus(TaskStatus.FINISHED_PENDING_CALLBACK);
taskMapper.updateById(task);
TaskEvent event = TaskEvent.builder()
.eventType("TASK_FINISHED")
.taskId(task.getId())
.taskNo(task.getTaskNo())
.ownerService(task.getOwnerService())
.bizType(task.getBizType())
.bizId(task.getBizId())
.handlerCode(task.getHandlerCode())
.payload(task.getRequestPayload())
.build();
mqProducer.send("TASK_EVENT_TOPIC", task.getOwnerService(), event);
}
}
```
注意Task 服务只负责发出“这个任务完成了”的事实事件,不直接调用 `StockAreaCallTubeTask`
## 6. LMS 如何精准定位到 `StockAreaCallTubeTask` 业务
LMS 消费到消息后,不是通过 Java 类名反射,而是通过 `handlerCode` 查本服务内注册的 handler。
### 6.1 定义统一 Handler 接口
```java
public interface TaskBizEventHandler {
String getHandlerCode();
void onTaskFinished(TaskEventMessage message);
void onTaskCancelled(TaskEventMessage message);
}
```
### 6.2 StockAreaCallTube 对应新 Handler
旧版类:
```text
StockAreaCallTubeTask
```
新版建议拆为:
```text
StockAreaCallTubeTaskCreator可选负责创建任务前业务校验和调用 Task 服务
StockAreaCallTubeTaskHandler负责完成/取消业务回调
```
例如:
```java
@Component
public class StockAreaCallTubeTaskHandler implements TaskBizEventHandler {
@Override
public String getHandlerCode() {
return "LMS_STOCK_AREA_CALL_TUBE";
}
@Override
@Transactional(rollbackFor = Exception.class)
public void onTaskFinished(TaskEventMessage message) {
// 旧 StockAreaCallTubeTask#updateTaskStatus FINISHED 的业务逻辑迁移到这里
// 1. 查询 LMS 业务记录或根据 taskId 查询业务关联
// 2. 更新终点备货位 ivt_status = 有货
// 3. 设置终点 vehicle_code
// 4. 清空起点 vehicle_code
// 5. 设置起点 ivt_status = 空
// 6. 根据 payload.to_material 恢复分切计划 is_paper_ok
// 7. 记录业务完成日志
}
@Override
@Transactional(rollbackFor = Exception.class)
public void onTaskCancelled(TaskEventMessage message) {
// 旧 StockAreaCallTubeTask#cancel 或 updateTaskStatus 取消逻辑迁移到这里
// 注意:旧版取消只是把任务状态置完成,新版应明确业务取消语义
}
}
```
### 6.3 LMS 内部 Handler 注册表
LMS 启动时把所有 handler 收集成 Map
```java
@Component
public class TaskBizEventHandlerRegistry {
private final Map<String, TaskBizEventHandler> handlerMap;
public TaskBizEventHandlerRegistry(List<TaskBizEventHandler> handlers) {
this.handlerMap = handlers.stream()
.collect(Collectors.toMap(TaskBizEventHandler::getHandlerCode, h -> h));
}
public TaskBizEventHandler getRequiredHandler(String handlerCode) {
TaskBizEventHandler handler = handlerMap.get(handlerCode);
if (handler == null) {
throw new IllegalArgumentException("未找到任务业务处理器:" + handlerCode);
}
return handler;
}
}
```
### 6.4 LMS 消费者分发
```java
@Component
public class LmsTaskEventConsumer {
@Autowired
private TaskBizEventHandlerRegistry registry;
public void onMessage(TaskEventMessage message) {
if (!"LMS".equals(message.getOwnerService())) {
return;
}
TaskBizEventHandler handler = registry.getRequiredHandler(message.getHandlerCode());
if ("TASK_FINISHED".equals(message.getEventType())) {
handler.onTaskFinished(message);
} else if ("TASK_CANCELLED".equals(message.getEventType())) {
handler.onTaskCancelled(message);
}
}
}
```
这样即使 LMS 里有很多任务处理器,例如:
```text
LMS_STOCK_AREA_CALL_TUBE
LMS_SEND_AIR_SHAFT_AGV
LMS_SLITTER_DOWN_AGV
LMS_TRUSS_CALL_SHAFT_CACHE
LMS_STOCK_AREA_SEND_VEHICLE
```
也能通过 `handlerCode` 精准定位到对应 handler。
## 7. 新旧机制对照
`StockAreaCallTubeTask` 为例:
| 步骤 | 旧版 | 新版 |
| --- | --- | --- |
| 创建任务 | LMS 内 `StockAreaCallTubeTask#createTask` | LMS 创建业务记录后调用 Task.createTask |
| 路由字段 | `handle_class = Java类名` | `owner_service = LMS`, `handler_code = LMS_STOCK_AREA_CALL_TUBE` |
| 下发 ACS | `AbstractAcsTask#immediateTaskNotifyAcs` | Task 服务 `AcsTaskDispatcher` |
| ACS 完成反馈 | LMS 接收 | Task 服务接收 |
| 找处理类 | `Class.forName(handle_class)` | LMS 消费 MQ 后用 `handlerCode` 查 Map |
| 完成逻辑 | `StockAreaCallTubeTask#updateTaskStatus` | `StockAreaCallTubeTaskHandler#onTaskFinished` |
| 取消逻辑 | `StockAreaCallTubeTask#cancel/updateTaskStatus` | `StockAreaCallTubeTaskHandler#onTaskCancelled` |
## 8. 消息中必须携带哪些字段
为了精准路由,消息至少包含:
```json
{
"eventId": "事件ID",
"eventType": "TASK_FINISHED",
"taskId": 10001,
"ownerService": "LMS",
"handlerCode": "LMS_STOCK_AREA_CALL_TUBE",
"bizType": "STOCK_AREA_CALL_TUBE",
"bizId": "业务ID",
"payload": {}
}
```
字段解释:
| 字段 | 是否必需 | 作用 |
| --- | --- | --- |
| `eventId` | 是 | 幂等消费 |
| `eventType` | 是 | 区分完成、取消、执行中 |
| `taskId` | 是 | Task 服务任务 ID |
| `ownerService` | 是 | 路由到 LMS/WMS |
| `handlerCode` | 是 | LMS/WMS 内部精准定位 handler |
| `bizType` | 建议 | 查询、统计、辅助判断 |
| `bizId` | 建议 | 业务 handler 回查业务数据 |
| `payload` | 建议 | 传递旧 request_param/result_param |
## 9. 业务处理成功后如何通知 Task
MQ 模式下LMS/WMS 处理完成后建议再通知 Task。
有两种方式。
### 9.1 方式一:业务服务 RPC 回调 Task
LMS 处理成功后调用:
```text
POST /rpc-api/task/callback-result
```
```json
{
"taskId": 10001,
"eventId": "EVT-10001",
"result": "SUCCESS",
"message": "StockAreaCallTube 完成业务处理成功"
}
```
Task 收到后:
```text
FINISHED_PENDING_CALLBACK -> FINISHED
```
### 9.2 方式二:业务服务发布处理结果 MQ
LMS 发布:
```text
TASK_CALLBACK_RESULT_TOPIC
```
Task 消费后更新最终状态。
第一阶段建议使用 RPC 回调 Task链路更直观。
## 10. 取消逻辑如何路由
取消和完成是同一套路由机制,只是事件类型不同。
Task 取消流程:
```text
用户取消任务
-> Task 判断状态
-> 如已下发则调用 ACS 取消
-> 外部取消成功
-> Task 状态改为 CANCEL_PENDING_CALLBACK
-> 发布 TASK_CANCELLED 事件Tag = LMS
-> LMS 消费消息
-> 根据 handlerCode = LMS_STOCK_AREA_CALL_TUBE 找 handler
-> 执行 onTaskCancelled
-> LMS 通知 Task 取消业务处理成功
-> Task 状态改为 CANCELLED
```
消息示例:
```json
{
"eventType": "TASK_CANCELLED",
"taskId": 10001,
"ownerService": "LMS",
"handlerCode": "LMS_STOCK_AREA_CALL_TUBE",
"bizType": "STOCK_AREA_CALL_TUBE",
"bizId": "LMS-BIZ-10001",
"payload": {
"cancelReason": "人工取消"
}
}
```
## 11. 为什么不是只靠 Topic 区分任务类
Topic 或 Tag 只能解决第一层路由:发给 LMS 还是 WMS。
它不能解决 LMS 内部多个任务类的问题。
例如 LMS 可能同时有:
```text
StockAreaCallTubeTask
StockAreaSendVehicleTask
SendAirShaftAgvTask
SlitterDownAgvTask
MoveVehicleAgvTask
```
如果只用 Topic/Tag只能知道这是一条 LMS 消息,不能知道该执行哪一个业务逻辑。
所以必须在消息体中携带:
```text
handlerCode
```
真正精准定位依赖的是:
```text
handlerCode -> LMS 内部 handlerMap -> 具体 Handler Bean
```
## 12. StockAreaCallTubeTask 迁移建议
旧版 `StockAreaCallTubeTask` 可以拆成两个类。
### 12.1 创建侧
```text
StockAreaCallTubeTaskCreator
```
职责:
- 校验备货区纸管搬运条件;
- 计算起点、终点、载具、product_area
- 生成 LMS 业务记录;
- 调用 Task.createTask
- 传入 `handlerCode = LMS_STOCK_AREA_CALL_TUBE`
### 12.2 回调侧
```text
StockAreaCallTubeTaskHandler
```
职责:
- 处理 `TASK_FINISHED`
- 处理 `TASK_CANCELLED`
- 迁移旧 `updateTaskStatus` 中的业务逻辑。
旧版完成逻辑迁移点:
```text
1. 更新终点备货位状态为有货
2. 将 vehicle_code 写入终点
3. 清空起点 vehicle_code
4. 起点状态改为空
5. 读取 to_material
6. 恢复对应分切计划 is_paper_ok
7. 更新业务日志
```
## 13. 幂等设计
### 13.1 Task 服务幂等
Task 接收 ACS 反馈时按以下维度幂等:
```text
taskId + externalEventId + status
```
同一个 ACS 完成反馈不能重复发布 `TASK_FINISHED`
### 13.2 LMS/WMS 消费幂等
业务服务按以下维度幂等:
```text
eventId
```
或:
```text
taskId + eventType
```
建议增加业务消费日志表:
```text
biz_task_event_consume_log
```
字段:
```text
event_id
task_id
handler_code
event_type
status
error_msg
```
如果已成功消费,再次收到直接返回成功。
## 14. 最终建议
新版本从 ACS 完成反馈到具体任务业务逻辑,不再靠跨服务反射 Java 类,而是两级路由:
```text
第一级owner_service / MQ Tag
决定消息给 LMS 还是 WMS
第二级handler_code
决定 LMS/WMS 内部哪个业务 handler 执行
```
`StockAreaCallTubeTask` 为例:
```text
ACS 完成反馈
-> Task 服务查询任务
-> owner_service = LMS
-> handler_code = LMS_STOCK_AREA_CALL_TUBE
-> 发布 TASK_FINISHEDTag = LMS
-> LMS 消费消息
-> handlerMap.get("LMS_STOCK_AREA_CALL_TUBE")
-> StockAreaCallTubeTaskHandler#onTaskFinished
```
这套设计既解决了 Task 服务不能依赖 LMS/WMS 的问题,又保留了旧版每种任务有自己业务处理类的扩展能力。