Files
huachuang/doc/task-service-technical-development.md

13 KiB
Raw Permalink Blame History

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 业务侧记录 IDLMS/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 业务回调是否成功