15 KiB
新版本任务完成/取消如何精准路由到具体业务 Handler
1. 问题背景
旧版兰州 LMS 中,任务类例如:
StockAreaCallTubeTask
备货区送纸管到机械手旁边的备货区 - AGV任务
它继承 AbstractAcsTask,创建任务时写入:
handle_class = org.nl.b_lms.sch.tasks.slitter.StockAreaCallTubeTask
当 ACS 请求任务完成时,旧系统会根据任务表里的 handle_class 找到对应 Spring Bean,然后执行:
StockAreaCallTubeTask#updateTaskStatus(taskObj, status)
所以旧版的路由链路是:
ACS 反馈完成
-> LMS 接收反馈
-> 查询 SCH_BASE_Task
-> 读取 handle_class
-> Class.forName(handle_class)
-> Spring 获取 StockAreaCallTubeTask Bean
-> 调用 updateTaskStatus
新版本中,任务模块独立为 Task 服务。Task 服务负责接收 ACS 反馈,但 StockAreaCallTubeTask 这类业务逻辑应该放在 LMS 服务中,而不是放在 Task 服务中。
因此新版本要解决两个路由问题:
- Task 服务如何知道这个任务完成后应该通知 LMS 还是 WMS?
- LMS/WMS 收到消息后,如何精准定位到原来类似
StockAreaCallTubeTask的业务逻辑?
2. 核心结论
新版本不再依赖 Java 类名 handle_class 跨服务反射,而是使用:
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 到业务服务 |
新版本完整链路:
ACS 反馈任务完成
-> Task 服务接收反馈
-> Task 查询任务记录
-> 得到 owner_service = LMS
-> 得到 handler_code = LMS_STOCK_AREA_CALL_TUBE
-> Task 发布 MQ:TaskFinishedEvent
-> LMS 消费属于自己的任务事件
-> LMS 根据 handler_code 找到 StockAreaCallTubeCallbackHandler
-> 执行 onTaskFinished
-> LMS 处理完成后回调 Task 或发送处理结果事件
-> Task 更新任务最终状态
3. 任务创建时必须保存路由信息
新版本能不能精准定位业务 handler,关键取决于创建任务时保存的信息是否完整。
以 StockAreaCallTubeTask 为例,旧版创建任务保存:
handle_class = org.nl.b_lms.sch.tasks.slitter.StockAreaCallTubeTask
新版本创建任务时,LMS 应调用 Task 服务并传入:
{
"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 |
业务侧记录 ID,handler 可用它查询业务数据 |
requestPayload |
保存旧任务中 request_param 类似的扩展数据 |
4. MQ Topic 应该如何设计
4.1 推荐方案:统一 Topic + owner_service 作为 Tag
推荐使用一个统一 Topic:
TASK_EVENT_TOPIC
不同业务服务通过 Tag 区分:
Tag = LMS
Tag = WMS
消息体里包含完整路由字段:
{
"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 只订阅:
Topic = TASK_EVENT_TOPIC
Tag = LMS
WMS 只订阅:
Topic = TASK_EVENT_TOPIC
Tag = WMS
优点:
- Topic 数量少;
- Task 事件模型统一;
- 新增业务服务只需要新增 Tag;
- 消息结构统一,便于日志、监控、重试。
4.2 备选方案:按服务拆 Topic
也可以拆成:
LMS_TASK_EVENT_TOPIC
WMS_TASK_EVENT_TOPIC
这种方式简单直接,但 Topic 会随着服务增多而增多。适合服务数量固定、消息隔离要求强的场景。
4.3 不推荐按任务类型拆 Topic
不推荐:
STOCK_AREA_CALL_TUBE_TOPIC
WMS_INBOUND_TOPIC
...
因为任务类型会越来越多,Topic 管理会失控。任务类型应该放在消息体的 bizType 或 handlerCode 中,而不是通过 Topic 区分。
5. Task 服务发布事件的逻辑
当 ACS 反馈任务完成时,Task 服务处理步骤:
1. 接收 ACS 反馈
2. 根据 taskId 查询 task_task
3. 幂等校验 external_event_id
4. 将任务状态改为 FINISHED_PENDING_CALLBACK
5. 根据 owner_service 决定消息 Tag
6. 发布 TASK_FINISHED 事件
7. 等待 LMS/WMS 业务处理结果
伪代码:
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 接口
public interface TaskBizEventHandler {
String getHandlerCode();
void onTaskFinished(TaskEventMessage message);
void onTaskCancelled(TaskEventMessage message);
}
6.2 StockAreaCallTube 对应新 Handler
旧版类:
StockAreaCallTubeTask
新版建议拆为:
StockAreaCallTubeTaskCreator,可选,负责创建任务前业务校验和调用 Task 服务
StockAreaCallTubeTaskHandler,负责完成/取消业务回调
例如:
@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:
@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 消费者分发
@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 里有很多任务处理器,例如:
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. 消息中必须携带哪些字段
为了精准路由,消息至少包含:
{
"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 处理成功后调用:
POST /rpc-api/task/callback-result
{
"taskId": 10001,
"eventId": "EVT-10001",
"result": "SUCCESS",
"message": "StockAreaCallTube 完成业务处理成功"
}
Task 收到后:
FINISHED_PENDING_CALLBACK -> FINISHED
9.2 方式二:业务服务发布处理结果 MQ
LMS 发布:
TASK_CALLBACK_RESULT_TOPIC
Task 消费后更新最终状态。
第一阶段建议使用 RPC 回调 Task,链路更直观。
10. 取消逻辑如何路由
取消和完成是同一套路由机制,只是事件类型不同。
Task 取消流程:
用户取消任务
-> Task 判断状态
-> 如已下发则调用 ACS 取消
-> 外部取消成功
-> Task 状态改为 CANCEL_PENDING_CALLBACK
-> 发布 TASK_CANCELLED 事件,Tag = LMS
-> LMS 消费消息
-> 根据 handlerCode = LMS_STOCK_AREA_CALL_TUBE 找 handler
-> 执行 onTaskCancelled
-> LMS 通知 Task 取消业务处理成功
-> Task 状态改为 CANCELLED
消息示例:
{
"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 可能同时有:
StockAreaCallTubeTask
StockAreaSendVehicleTask
SendAirShaftAgvTask
SlitterDownAgvTask
MoveVehicleAgvTask
如果只用 Topic/Tag,只能知道这是一条 LMS 消息,不能知道该执行哪一个业务逻辑。
所以必须在消息体中携带:
handlerCode
真正精准定位依赖的是:
handlerCode -> LMS 内部 handlerMap -> 具体 Handler Bean
12. StockAreaCallTubeTask 迁移建议
旧版 StockAreaCallTubeTask 可以拆成两个类。
12.1 创建侧
StockAreaCallTubeTaskCreator
职责:
- 校验备货区纸管搬运条件;
- 计算起点、终点、载具、product_area;
- 生成 LMS 业务记录;
- 调用 Task.createTask;
- 传入
handlerCode = LMS_STOCK_AREA_CALL_TUBE。
12.2 回调侧
StockAreaCallTubeTaskHandler
职责:
- 处理
TASK_FINISHED; - 处理
TASK_CANCELLED; - 迁移旧
updateTaskStatus中的业务逻辑。
旧版完成逻辑迁移点:
1. 更新终点备货位状态为有货
2. 将 vehicle_code 写入终点
3. 清空起点 vehicle_code
4. 起点状态改为空
5. 读取 to_material
6. 恢复对应分切计划 is_paper_ok
7. 更新业务日志
13. 幂等设计
13.1 Task 服务幂等
Task 接收 ACS 反馈时按以下维度幂等:
taskId + externalEventId + status
同一个 ACS 完成反馈不能重复发布 TASK_FINISHED。
13.2 LMS/WMS 消费幂等
业务服务按以下维度幂等:
eventId
或:
taskId + eventType
建议增加业务消费日志表:
biz_task_event_consume_log
字段:
event_id
task_id
handler_code
event_type
status
error_msg
如果已成功消费,再次收到直接返回成功。
14. 最终建议
新版本从 ACS 完成反馈到具体任务业务逻辑,不再靠跨服务反射 Java 类,而是两级路由:
第一级:owner_service / MQ Tag
决定消息给 LMS 还是 WMS
第二级:handler_code
决定 LMS/WMS 内部哪个业务 handler 执行
以 StockAreaCallTubeTask 为例:
ACS 完成反馈
-> Task 服务查询任务
-> owner_service = LMS
-> handler_code = LMS_STOCK_AREA_CALL_TUBE
-> 发布 TASK_FINISHED,Tag = LMS
-> LMS 消费消息
-> handlerMap.get("LMS_STOCK_AREA_CALL_TUBE")
-> StockAreaCallTubeTaskHandler#onTaskFinished
这套设计既解决了 Task 服务不能依赖 LMS/WMS 的问题,又保留了旧版每种任务有自己业务处理类的扩展能力。