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

631 lines
15 KiB
Markdown
Raw Normal View History

2026-07-14 14:09:48 +08:00
# 新版本任务完成/取消如何精准路由到具体业务 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 的问题,又保留了旧版每种任务有自己业务处理类的扩展能力。