13 KiB
13 KiB
Task 服务技术开发文档
1. 服务定位
Task 服务基于 task_transport_job 表,负责搬运任务的通用生命周期管理和 ACS 交互。
整体链路:
LMS/WMS 创建搬运任务
-> Task 保存任务
-> Task 自动或手动下发 ACS
-> ACS 反馈执行中、取货、完成、取消
-> Task 更新任务状态
-> Task 通过 MQ 通知 LMS/WMS 执行业务完成/取消
-> LMS/WMS 处理后回传结果
-> Task 更新最终完成/取消状态
Task 服务不处理 LMS/WMS 的具体业务数据,例如生产计划、备货区点位、库存、入库单、出库单。业务逻辑由 LMS/WMS 内部 handler 处理。
2. Task 服务除 CRUD 外的核心业务
Task 服务核心能力包括:
1. 接收 LMS/WMS 创建任务
2. 保存任务标准数据
3. 定时扫描并自动下发任务
4. 将 task_transport_job 转换为 AcsTaskDto
5. 调用 ACS 下发任务
6. 接收 ACS 状态反馈
7. 维护任务状态机
8. 发布完成/取消 MQ 给 LMS/WMS
9. 接收 LMS/WMS 业务处理结果
10. 处理下发失败、回调失败、超时补偿
11. 幂等和并发控制
12. 任务过程日志记录
3. 表字段和核心职责对应关系
3.1 业务路由字段
| 字段 | 作用 |
|---|---|
owner_service |
决定 MQ 发送给 LMS 还是 WMS |
biz_type |
业务类型,用于查询、统计、防重复创建 |
biz_id |
业务侧记录 ID,LMS/WMS 回调时定位业务数据 |
handle_code |
LMS/WMS 内部 handler 编码 |
完成/取消时路由:
owner_service -> MQ Tag 或目标服务
handle_code -> LMS/WMS 内部 handlerMap
3.2 ACS 下发字段
| 字段 | 对应 AcsTaskDto |
|---|---|
task_id |
taskId |
task_code |
taskCode |
acs_task_type |
taskType |
point_code1 |
startDeviceCode |
point_code2 |
nextDeviceCode |
point_code3 |
startDeviceCode2 |
point_code4 |
nextDeviceCode2 |
vehicle_code |
vehicleCode |
vehicle_code2 |
vehicleCode2 |
priority |
priority |
product_area |
productArea |
agv_system_type |
agvSystemType |
request_param |
interactionJson |
3.3 参数快照字段
| 字段 | 说明 |
|---|---|
request_param |
LMS/WMS 创建任务扩展参数,同时作为 interactionJson 来源 |
dispatch_param |
Task 下发 ACS 的 AcsTaskDto 报文快照 |
result_param |
ACS 反馈报文 |
3.4 回调字段
| 字段 | 说明 |
|---|---|
callback_status |
业务回调状态:PENDING/SUCCESS/FAILED |
callback_retry_count |
业务回调重试次数 |
callback_error_msg |
业务回调失败原因 |
4. 任务状态设计
推荐状态码:
| 状态码 | 状态名称 | 修改阶段 |
|---|---|---|
10 |
生成 | 创建任务但条件不完整 |
40 |
待下发 | 起点终点确认,允许下发 |
45 |
下发中 | Task 准备调用 ACS 前 |
50 |
已下发 | ACS 接收任务成功 |
60 |
执行中 | ACS 反馈执行中 |
61 |
已取货搬运中 | ACS 反馈取货完成,部分任务有 |
67 |
外部完成待业务处理 | ACS 完成,等待 LMS/WMS 业务完成 |
68 |
外部取消中 | Task 正在调用 ACS 取消 |
69 |
外部取消待业务处理 | ACS 取消成功,等待 LMS/WMS 回滚 |
70 |
完成 | LMS/WMS 业务完成成功 |
80 |
取消 | LMS/WMS 业务取消成功 |
90 |
失败 | 下发失败、取消失败、不可恢复异常 |
常用查询:
-- 外部未完成任务,用于防止重复创建
SELECT *
FROM task_transport_job
WHERE deleted = b'0'
AND biz_type = ?
AND task_status < '067';
-- 未最终结束任务
SELECT *
FROM task_transport_job
WHERE deleted = b'0'
AND task_status < '070';
-- 业务回调失败任务
SELECT *
FROM task_transport_job
WHERE deleted = b'0'
AND task_status IN ('067', '069')
AND callback_status = 'FAILED';
5. 创建任务
5.1 接口
POST /rpc-api/task/transport/create
5.2 创建请求核心字段
public class TransportTaskCreateReqDTO {
private String taskName;
private String ownerService;
private String bizType;
private String bizId;
private String handleCode;
private String taskType;
private String acsTaskType;
private String agvSystemType;
private String pointCode1;
private String pointCode2;
private String pointCode3;
private String pointCode4;
private String vehicleCode;
private String vehicleCode2;
private String priority;
private String productArea;
private String isAutoIssue;
private Map<String, Object> requestParam;
}
5.3 创建逻辑
1. 校验 owner_service、handle_code、acs_task_type 等必要字段
2. 生成 task_id、task_code
3. 保存业务归属字段 owner_service、biz_type、biz_id、handle_code
4. 保存 ACS 下发字段 point_code、vehicle_code、acs_task_type 等
5. 保存 request_param
6. 参数完整则 task_status = 040,否则 task_status = 010
6. 自动下发任务
6.1 定时器
推荐定时任务:
autoIssueTransportTaskJob
扫描条件:
SELECT *
FROM task_transport_job
WHERE deleted = b'0'
AND is_auto_issue = '1'
AND task_status = '040'
ORDER BY priority DESC, create_time ASC
LIMIT ?;
6.2 抢占任务
为避免并发重复下发,先做状态抢占:
UPDATE task_transport_job
SET task_status = '045', update_time = NOW()
WHERE task_id = ?
AND task_status = '040';
影响行数为 1 才继续下发。
6.3 下发流程
1. 查询 040 待下发任务
2. 抢占为 045 下发中
3. 转换成 AcsTaskDto
4. 保存 dispatch_param
5. 调用 ACS 下发接口
6. 成功:task_status = 050,保存 external_task_no
7. 失败:task_status = 090 或恢复 040,记录错误
7. AcsTaskDto 转换
7.1 DTO
public class AcsTaskDto {
private String taskId;
private String taskCode;
private String taskType;
private String startDeviceCode;
private String nextDeviceCode;
private String startDeviceCode2;
private String nextDeviceCode2;
private String vehicleCode;
private String vehicleCode2;
private String priority;
private String productArea;
private String agvSystemType;
private Map<String, Object> interactionJson;
}
7.2 转换器
@Component
public class AcsTaskDtoConverter {
public AcsTaskDto convert(TaskTransportJobDO task) {
AcsTaskDto dto = new AcsTaskDto();
dto.setTaskId(String.valueOf(task.getTaskId()));
dto.setTaskCode(task.getTaskCode());
dto.setTaskType(task.getAcsTaskType());
dto.setStartDeviceCode(task.getPointCode1());
dto.setNextDeviceCode(task.getPointCode2());
dto.setStartDeviceCode2(task.getPointCode3());
dto.setNextDeviceCode2(task.getPointCode4());
dto.setVehicleCode(task.getVehicleCode());
dto.setVehicleCode2(task.getVehicleCode2());
dto.setPriority(task.getPriority());
dto.setProductArea(task.getProductArea());
dto.setAgvSystemType(task.getAgvSystemType());
dto.setInteractionJson(JsonUtils.parseMap(task.getRequestParam()));
return dto;
}
}
8. ACS 下发服务
8.1 ACS Client
public interface AcsClient {
AcsIssueResult issueTask(AcsTaskDto taskDto);
AcsCancelResult cancelTask(String externalTaskNo, Long taskId);
}
8.2 下发服务职责
TransportTaskIssueService.issue(taskId)
-> 查询任务
-> 抢占状态 040 -> 045
-> AcsTaskDtoConverter.convert
-> 保存 dispatch_param
-> acsClient.issueTask
-> 成功 045 -> 050
-> 失败记录错误
9. ACS 反馈接收
9.1 接口
POST /api/task/transport/acs-feedback
9.2 反馈 DTO
public class AcsFeedbackReqDTO {
private Long taskId;
private String taskCode;
private String externalTaskNo;
private String status;
private String eventId;
private Map<String, Object> payload;
}
9.3 状态处理
| ACS 反馈 | Task 更新 |
|---|---|
| 执行中 | task_status = 060 |
| 已取货 | task_status = 061 |
| 完成 | task_status = 067, callback_status = PENDING |
| 取消成功 | task_status = 069, callback_status = PENDING |
所有反馈均保存到:
result_param
完成和取消反馈后,需要发布 MQ 给 LMS/WMS。
10. 完成/取消 MQ 事件
10.1 Topic 设计
Topic = TASK_EVENT_TOPIC
Tag = owner_service
示例:
TASK_EVENT_TOPIC:LMS
TASK_EVENT_TOPIC:WMS
10.2 消息结构
public class TaskEventMessage {
private String eventId;
private String eventType; // TASK_FINISHED / TASK_CANCELLED
private Long taskId;
private String taskCode;
private String ownerService;
private String bizType;
private String bizId;
private String handleCode;
private Map<String, Object> payload;
}
10.3 发布逻辑
1. 查询 task_transport_job
2. 构造 TaskEventMessage
3. topic = TASK_EVENT_TOPIC
4. tag = owner_service
5. 发送 MQ
6. 发送失败:callback_status = FAILED,记录 callback_error_msg
11. 接收业务回调结果
LMS/WMS 消费 MQ 后,内部通过 handle_code 找 handler 执行业务逻辑。处理完成后回调 Task。
接口:
POST /rpc-api/task/transport/callback-result
DTO:
public class TaskCallbackResultReqDTO {
private Long taskId;
private String eventId;
private String eventType;
private String result; // SUCCESS / FAILED
private String message;
}
状态更新:
| 场景 | Task 更新 |
|---|---|
| 完成业务成功 | task_status = 070, callback_status = SUCCESS |
| 完成业务失败 | task_status = 067, callback_status = FAILED, callback_retry_count + 1 |
| 取消业务成功 | task_status = 080, callback_status = SUCCESS |
| 取消业务失败 | task_status = 069, callback_status = FAILED, callback_retry_count + 1 |
12. 人工完成和取消
12.1 人工完成
PC 点击完成
-> Task 校验任务
-> task_status = 067
-> callback_status = PENDING
-> 发布 TASK_FINISHED MQ
-> 页面提示:完成请求已提交,业务处理中
12.2 人工取消
未下发任务:
010/040 -> 069 -> 发布 TASK_CANCELLED MQ
已下发任务:
050/060/061 -> 068 -> 调 ACS 取消 -> 069 -> 发布 TASK_CANCELLED MQ
13. 补偿和重试
13.1 回调失败重试
扫描:
SELECT *
FROM task_transport_job
WHERE deleted = b'0'
AND task_status IN ('067', '069')
AND callback_status = 'FAILED'
AND callback_retry_count < ?;
处理:
重新发布 TASK_FINISHED 或 TASK_CANCELLED
callback_retry_count + 1
13.2 PENDING 超时补偿
扫描:
task_status IN ('067', '069')
callback_status = PENDING
update_time 超过阈值
处理:
重新发布 MQ 或标记 FAILED
13.3 下发中超时补偿
扫描:
task_status = 045
update_time 超过阈值
处理:
查询 ACS 是否已收到任务;无法确认则恢复 040 或标记 090
14. 幂等控制
14.1 创建幂等
防重复创建示例:
SELECT *
FROM task_transport_job
WHERE deleted = b'0'
AND owner_service = ?
AND biz_type = ?
AND biz_id = ?
AND task_status < '067';
14.2 下发幂等
通过 040 -> 045 条件更新抢占。
14.3 ACS 反馈幂等
建议引入反馈日志表记录 eventId。第一版至少根据状态判断:
已是 070,不重复处理完成
已是 067,不重复发布完成 MQ
14.4 业务回调结果幂等
如果已是最终态:
task_status = 070 或 080
再次收到成功回调时直接返回成功。
15. 推荐模块结构
nl-module-task-server
controller
TransportTaskController
AcsFeedbackController
api
TransportTaskApiImpl
service
TransportTaskService
TransportTaskIssueService
TransportTaskFeedbackService
TransportTaskCallbackService
TransportTaskCancelService
TransportTaskRetryService
converter
AcsTaskDtoConverter
mq
TaskEventProducer
job
AutoIssueTransportTaskJob
CallbackRetryJob
dal
dataobject/TaskTransportJobDO
mapper/TaskTransportJobMapper
enums
TransportTaskStatusEnum
CallbackStatusEnum
TaskEventTypeEnum
16. 总结
Task 服务的核心不是简单 CRUD,而是搬运任务生命周期中台。
核心职责:
接收任务 -> 保存任务 -> 自动下发 -> 转换 AcsTaskDto -> 调 ACS -> 接收反馈 -> 发布业务事件 -> 接收业务结果 -> 最终完成/取消 -> 补偿重试
关键路由:
owner_service 决定事件给 LMS/WMS
handle_code 决定业务服务内部哪个 handler 处理
关键状态:
task_status 表示任务流转到哪一步
callback_status 表示 LMS/WMS 业务回调是否成功