From d9f21147aafaf6d42f8431ead0f794030fc7d541 Mon Sep 17 00:00:00 2001 From: liyongde <1419499670@qq.com> Date: Mon, 3 Aug 2026 19:30:54 +0800 Subject: [PATCH] =?UTF-8?q?feat=EF=BC=9A=E5=8F=96=E6=B6=88=E5=AE=8C?= =?UTF-8?q?=E6=88=90=E9=80=9A=E8=BF=87=E9=85=8D=E7=BD=AE=E5=86=B3=E5=AE=9A?= =?UTF-8?q?=E6=98=AF=E5=90=A6=E4=BD=BF=E7=94=A8mq=E6=88=96=E8=80=85feign?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../plans/2026-08-03-task-callback-mode.md | 67 ++++++++ .../2026-08-03-task-callback-mode-design.md | 72 +++++++++ .../nl-spring-boot-starter-execute/pom.xml | 9 +- .../biz/api/AbstractTaskCommonApiImpl.java | 25 +++ .../execute/biz/api/TaskCommonApi.java | 8 + .../biz/dto/TaskStatusCallApiReqDTO.java | 3 + .../execute/core/TaskCommonApiFactory.java | 8 + .../execute/core/dto/TaskExecuteDTO.java | 6 + .../framework/execute/util/http/AcsUtil.java | 2 +- ...tractTaskCommonApiImplLocalSpringTest.java | 126 +++++++++++++++ .../TaskCommonApiFactoryLocalSpringTest.java | 61 +++++++ .../consumer/LmsTaskStatusChangeConsumer.java | 57 ++++++- nl-module-task/nl-module-task-api/pom.xml | 9 +- .../task/enums/TaskCallbackTypeEnum.java | 35 ++++ .../task/enums/TaskOwnerServiceEnum.java | 50 ++++++ .../module/task/message/TaskEventMessage.java | 11 ++ .../TaskOwnerServiceEnumLocalSpringTest.java | 56 +++++++ .../manage/TransportTaskOperationManager.java | 152 ++++++++++++++---- .../task/mq/producer/TaskEventProducer.java | 11 +- .../TransportTaskCallbackServiceImpl.java | 61 ++++--- .../src/main/resources/application-dev.yaml | 20 +-- .../src/main/resources/application.yaml | 2 + .../consumer/WmsTaskStatusChangeConsumer.java | 57 ++++++- .../src/main/resources/application-dev.yaml | 8 +- nl-server/src/main/resources/application.yaml | 11 +- .../nl/server/acs/AcsUtilLocalSpringTest.java | 2 +- 26 files changed, 841 insertions(+), 88 deletions(-) create mode 100644 docs/superpowers/plans/2026-08-03-task-callback-mode.md create mode 100644 docs/superpowers/specs/2026-08-03-task-callback-mode-design.md create mode 100644 nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImplLocalSpringTest.java create mode 100644 nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/core/TaskCommonApiFactoryLocalSpringTest.java create mode 100644 nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskCallbackTypeEnum.java create mode 100644 nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskOwnerServiceEnum.java create mode 100644 nl-module-task/nl-module-task-api/src/test/java/cn/code/nl/module/task/enums/TaskOwnerServiceEnumLocalSpringTest.java diff --git a/docs/superpowers/plans/2026-08-03-task-callback-mode.md b/docs/superpowers/plans/2026-08-03-task-callback-mode.md new file mode 100644 index 00000000..a105c319 --- /dev/null +++ b/docs/superpowers/plans/2026-08-03-task-callback-mode.md @@ -0,0 +1,67 @@ +# 任务完成与取消双通道回调实现计划 + +> **面向 AI 代理的工作者:** 在当前会话逐任务实现;步骤使用复选框跟踪。当前工作区存在用户改动,不自动提交、不覆盖无关文件。 + +**目标:** 实现 YAML 驱动的 MQ 异步与 Feign 同步完成/取消链路,并在业务成功后统一推进任务最终状态。 + +**架构:** Task 操作管理器负责配置分流,execute starter 提供统一 Feign 业务入口,MQ Consumer 在业务处理后回调 Task。最终状态仅由 `TransportTaskCallbackService` 收口更新。 + +**技术栈:** Java 17、Spring Boot、OpenFeign、RocketMQ、MyBatis、JUnit 5 + +--- + +### 任务 1:补齐 execute starter 的完成与取消入口 + +**文件:** +- 修改:`nl-framework/nl-spring-boot-starter-execute/pom.xml` +- 修改:`nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java` +- 修改:`nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImpl.java` +- 创建:`nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImplLocalSpringTest.java` + +- [x] 编写 Spring 集成测试,调用 `doHandleFinished`、`doHandleCancelled` 并断言测试任务处理器收到正确 taskId 与 payload。 +- [x] 运行定向测试,确认因接口缺失而失败。 +- [x] 在 `TaskCommonApi` 和公共实现中增加两个接口,复用现有 `execute` 路由。 +- [x] 重跑定向测试并确认通过。 + +### 任务 2:增加 YAML 模式分流并统一同步状态收口 + +**文件:** +- 创建:`nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskCallbackTypeEnum.java` +- 修改:`nl-module-task/nl-module-task-server/src/main/resources/application.yaml` +- 修改:`nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java` +- 修改:`nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java` + +- [x] 定义 `mq`、`feign` 配置枚举及严格解析方法。 +- [x] 为 Task 服务增加 `nl.task.callback-type: mq`。 +- [x] 将完成、取消共同的待处理状态保存后,按配置选择事务后发 MQ 或同步调用 `TaskCommonApi`。 +- [x] Feign 成功后构造成功回调请求并调用 `TransportTaskCallbackService`。 +- [x] 删除字符串状态码的错误格式化,直接使用枚举 code。 + +### 任务 3:补齐 MQ 成功与失败回调 + +**文件:** +- 修改:`nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/consumer/LmsTaskStatusChangeConsumer.java` +- 修改:`nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java` + +- [x] 完成或取消业务正常结束后,通过 `TransportTaskApi.receiveCallbackResult` 回调 `SUCCESS` 并检查返回结果。 +- [x] 业务异常时回调 `FAILED` 记录异常信息,再抛出原异常触发 MQ 重试。 +- [x] 未找到处理器时抛出异常,避免消息被确认后任务永久停留在待处理状态。 +- [x] 统一事件标识并使用 Redis 成功标识隔离业务执行与成功回调重试。 + +### 任务 4:验证与审查 + +**文件:** +- 检查:以上全部生产与测试文件 + +- [ ] 运行 execute starter 定向 Spring 集成测试。 +- [ ] 编译 execute starter、task-api、task-server、lms-server、wms-server 及依赖模块。 +- [ ] 检查 git diff,确认没有覆盖工作区既有改动。 +- [ ] 按 Java 规范检查中文注释、`@Resource` 注入、Feign 返回值检查及异常日志。 + +### 最终验证结果 + +- [x] execute starter 4 个 Spring 定向测试通过。 +- [x] task-api 3 个业务归属与 Topic 映射 Spring 定向测试通过。 +- [x] Execute、Task、LMS、WMS 及依赖共 33 个模块编译通过。 +- [x] 聚合启动工程 `nl-server` 及依赖共 38 个模块编译通过。 +- [x] `git diff --check` 通过,仅有工作区行尾转换提示。 diff --git a/docs/superpowers/specs/2026-08-03-task-callback-mode-design.md b/docs/superpowers/specs/2026-08-03-task-callback-mode-design.md new file mode 100644 index 00000000..52b840cf --- /dev/null +++ b/docs/superpowers/specs/2026-08-03-task-callback-mode-design.md @@ -0,0 +1,72 @@ +# 任务完成与取消双通道回调设计 + +## 目标 + +任务模块根据 YAML 配置选择 MQ 异步或 Feign 同步方式通知 LMS/WMS 执行完成、取消业务;业务处理成功后,两条链路都由 Task 服务统一把任务从回调待处理状态推进到最终完成或取消状态。 + +## 当前问题 + +- `TransportTaskOperationManager#handleFinished`、`handleCancelled` 固定发布 RocketMQ 消息,无法切换为同步 Feign。 +- LMS/WMS 的 MQ Consumer 调用 `AbstractTask#doHandleFinish`、`doHandleCancel` 后没有回调 Task,任务停留在 `FINISHED_CALLBACK_PENDING` 或 `CANCEL_CALLBACK_PENDING`。 +- `TransportTaskCallbackServiceImpl` 使用 `%03d` 格式化字符串状态码,回调处理存在运行时类型异常。 +- `TaskCommonApi` 只暴露取货、进出等动作,没有完成和取消接口。 + +## 配置 + +在 Task 服务 YAML 中增加必填配置: + +```yaml +nl: + task: + callback-type: mq +``` + +允许值为 `mq`、`feign`。使用枚举解析配置;配置缺失或值非法时启动失败,不做静默兜底。默认项目配置使用 `mq`,保持现有部署行为。 + +## 执行链路 + +### MQ 异步模式 + +1. Task 将任务置为对应的回调待处理状态,并把 `callbackStatus` 置为 `PENDING`。 +2. 当前事务提交后发布完成或取消消息。 +3. LMS/WMS Consumer 校验任务仍处于当前事件对应的待处理状态,然后路由到 `AbstractTask` 子类执行业务。 +4. 业务成功后,Consumer 写入带有效期的 Redis 成功标识,再通过 `TransportTaskApi.receiveCallbackResult` 回调 `SUCCESS`;若成功回调失败,MQ 重试时只重试回调,不重复执行业务。 +5. Task 的统一回调服务把任务更新为 `FINISHED` 或 `CANCELLED`,并把 `callbackStatus` 置为 `SUCCESS`。 +6. 业务异常时,Consumer 尝试回调 `FAILED` 记录错误与重试次数,再重新抛出异常,让 MQ 执行重试。 + +### Feign 同步模式 + +1. Task 将任务置为对应的回调待处理状态,并把 `callbackStatus` 置为 `PENDING`。 +2. Task 根据 `ownerService` 从 `TaskCommonApiFactory` 获取 LMS/WMS Feign Client。 +3. 调用新增的 `doHandleFinished` 或 `doHandleCancelled`,业务服务通过 `TaskFactory` 路由到相同的 `AbstractTask` 子类。 +4. Feign 返回成功后,Task 直接调用统一回调服务,以相同规则更新最终状态。 +5. Feign 返回失败或抛出异常时不进入最终状态,异常沿现有任务操作链路上抛。 + +## 接口与数据 + +- 在 `TaskCommonApi` 增加 `doHandleFinished`、`doHandleCancelled`。 +- 两个接口继续复用 `TaskStatusCallApiReqDTO`,携带任务标识、业务归属、业务标识、处理器编码和 payload。 +- 同步与异步链路都复用 `TaskCallbackResultReqDTO` 和 `TransportTaskCallbackService` 完成最终状态收口。 +- `eventId` 使用任务 ID 与事件类型组成的稳定值;Task 按当前 pending 状态校验 eventId,拒绝延迟或错序事件推进错误终态。 + +## 幂等与错误处理 + +- Task 已处于 `FINISHED` 或 `CANCELLED` 时忽略重复成功回调。 +- 只有 `FINISHED_CALLBACK_PENDING` 能推进到 `FINISHED`,只有 `CANCEL_CALLBACK_PENDING` 能推进到 `CANCELLED`。 +- MQ 重复消息在消费前查询 Task 状态;状态已不匹配时直接跳过。 +- MQ 业务成功后使用七天 Redis 标识隔离业务执行与 Task 成功回调重试;成功回调完成后删除标识。 +- Feign `CommonResult` 必须调用 `getCheckedData` 或检查 `isSuccess`,业务失败不允许误更新最终状态。 +- MQ 回调失败不吞掉业务异常,保留 RocketMQ 重试能力。 +- 任务状态分布式锁延迟到本地事务完成后释放,避免提交前窗口发生并发重复处理。 + +## 验证 + +- Spring 集成测试验证 execute starter 的完成、取消 Feign 入口能正确路由到任务处理器。 +- 编译 execute starter、task-api、task-server、lms-server、wms-server,验证跨模块接口一致。 +- 检查 MQ 成功回调和 Feign 成功返回最终都进入统一状态处理服务。 + +## 一致性边界 + +- MQ 和 Feign 均向 `AbstractTask` 透传稳定的 `eventId` 与 `eventType`,业务处理器应以该事件标识在业务数据库中实现幂等。 +- Redis 成功标识可以避免“业务已成功、Task 成功回调失败”时重复执行业务,但不能覆盖业务事务提交后、Redis 标识写入前进程崩溃的极小窗口,因此业务数据库幂等仍是最终保障。 +- Feign 调用与 Task 本地事务不构成分布式事务;远端成功后若本地事务回滚,重试仍依赖上述事件幂等。 diff --git a/nl-framework/nl-spring-boot-starter-execute/pom.xml b/nl-framework/nl-spring-boot-starter-execute/pom.xml index 093829fa..8eccfdec 100644 --- a/nl-framework/nl-spring-boot-starter-execute/pom.xml +++ b/nl-framework/nl-spring-boot-starter-execute/pom.xml @@ -47,6 +47,13 @@ springdoc-openapi-starter-webmvc-ui provided + + + + org.springframework.boot + spring-boot-starter-test + test + - \ No newline at end of file + diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImpl.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImpl.java index 9f4fb1e5..5a331dd5 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImpl.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImpl.java @@ -22,6 +22,29 @@ public abstract class AbstractTaskCommonApiImpl implements TaskCommonApi { @Resource private TaskFactory taskFactory; + /** + * 处理任务完成 + */ + @Override + public CommonResult doHandleFinished(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { + return execute(taskStatusCallApiReqDTO, (task, taskExecuteDTO) -> { + task.doHandleFinish(taskExecuteDTO); + + return null; + }); + } + + /** + * 处理任务取消 + */ + @Override + public CommonResult doHandleCancelled(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { + return execute(taskStatusCallApiReqDTO, (task, taskExecuteDTO) -> { + task.doHandleCancel(taskExecuteDTO); + return null; + }); + } + /** * 处理取货完成 */ @@ -97,6 +120,8 @@ public abstract class AbstractTaskCommonApiImpl implements TaskCommonApi { private TaskExecuteDTO buildTaskExecuteDTO(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { return TaskExecuteDTO.builder() .taskId(taskStatusCallApiReqDTO.getTaskId()) + .eventId(taskStatusCallApiReqDTO.getEventId()) + .eventType(taskStatusCallApiReqDTO.getStatus()) .payload(taskStatusCallApiReqDTO.getPayload()) .build(); } diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java index fff86e58..df462d9a 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/api/TaskCommonApi.java @@ -14,6 +14,14 @@ import org.springframework.web.bind.annotation.RequestBody; */ public interface TaskCommonApi { + @Operation(summary = "处理任务完成") + @PostMapping("/doHandleFinished") + CommonResult doHandleFinished(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); + + @Operation(summary = "处理任务取消") + @PostMapping("/doHandleCancelled") + CommonResult doHandleCancelled(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); + @Operation(summary = "请求取货") @PostMapping("/do-handle-picked") CommonResult doHandlePicked(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/dto/TaskStatusCallApiReqDTO.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/dto/TaskStatusCallApiReqDTO.java index b3085140..f2a8e175 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/dto/TaskStatusCallApiReqDTO.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/biz/dto/TaskStatusCallApiReqDTO.java @@ -15,6 +15,9 @@ public class TaskStatusCallApiReqDTO implements Serializable { /** 任务ID */ private Long taskId; + /** 任务事件唯一标识 */ + private String eventId; + /** 任务编码 */ private String taskCode; diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/TaskCommonApiFactory.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/TaskCommonApiFactory.java index 7a620b26..2a1587f8 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/TaskCommonApiFactory.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/TaskCommonApiFactory.java @@ -23,6 +23,12 @@ import java.util.Map; @Component public class TaskCommonApiFactory { + /** LMS 业务归属编码 */ + private static final String LMS_OWNER_SERVICE = "LMS"; + + /** WMS 业务归属编码 */ + private static final String WMS_OWNER_SERVICE = "WMS"; + @Autowired(required = false) private LmsTaskCommonApi lmsTaskCommonApi; @@ -41,9 +47,11 @@ public class TaskCommonApiFactory { public void init() { if (lmsTaskCommonApi != null) { apiMap.put(RpcConstants.LMS_NAME, lmsTaskCommonApi); + apiMap.put(LMS_OWNER_SERVICE, lmsTaskCommonApi); } if (wmsTaskCommonApi != null) { apiMap.put(RpcConstants.WMS_NAME, wmsTaskCommonApi); + apiMap.put(WMS_OWNER_SERVICE, wmsTaskCommonApi); } } diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/dto/TaskExecuteDTO.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/dto/TaskExecuteDTO.java index e0cff2a4..42e9a813 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/dto/TaskExecuteDTO.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/core/dto/TaskExecuteDTO.java @@ -24,6 +24,12 @@ public class TaskExecuteDTO { @NotNull(message = "任务ID不能为空") private Long taskId; + @Schema(description = "任务事件唯一标识") + private String eventId; + + @Schema(description = "任务事件类型") + private String eventType; + @Schema(description = "ACS反馈扩展数据") private Map payload; } diff --git a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/util/http/AcsUtil.java b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/util/http/AcsUtil.java index 81fd560d..a66ee4f7 100644 --- a/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/util/http/AcsUtil.java +++ b/nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/util/http/AcsUtil.java @@ -53,7 +53,7 @@ public class AcsUtil { String response; try { response = HttpUtils.post(url, headers(), JsonUtils.toJsonString(request)); - log.info("ACS 返回值:{}", JSON.toJSONString(response)); + log.info("ACS 返回值:{}", response); } catch (RuntimeException ex) { if (isNetworkException(ex)) { log.error("ACS 服务网络不通,url={},request={}", url, JsonUtils.toJsonString(request), ex); diff --git a/nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImplLocalSpringTest.java b/nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImplLocalSpringTest.java new file mode 100644 index 00000000..b197b5e5 --- /dev/null +++ b/nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/biz/api/AbstractTaskCommonApiImplLocalSpringTest.java @@ -0,0 +1,126 @@ +package cn.code.nl.framework.execute.biz.api; + +import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO; +import cn.code.nl.framework.execute.core.AbstractTask; +import cn.code.nl.framework.execute.core.TaskFactory; +import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO; +import jakarta.annotation.Resource; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.stereotype.Component; + +import java.util.Map; + +/** + * 任务通用 Feign 接口本地 Spring 集成测试 + */ +@SpringBootTest(classes = { + TaskFactory.class, + AbstractTaskCommonApiImplLocalSpringTest.TestTaskCommonApi.class, + AbstractTaskCommonApiImplLocalSpringTest.TestTask.class +}) +class AbstractTaskCommonApiImplLocalSpringTest { + + @Resource + private TestTaskCommonApi testTaskCommonApi; + + @Resource + private TestTask testTask; + + /** + * 验证完成接口按 handleCode 路由并传递任务参数 + */ + @Test + void doHandleFinishedShouldRouteTask() { + TaskStatusCallApiReqDTO reqDTO = buildReqDTO(); + + testTaskCommonApi.doHandleFinished(reqDTO).getCheckedData(); + + Assertions.assertNotNull(testTask.getFinishedDTO()); + Assertions.assertEquals(1001L, testTask.getFinishedDTO().getTaskId()); + Assertions.assertEquals("1001_TASK_FINISHED", testTask.getFinishedDTO().getEventId()); + Assertions.assertEquals("TASK_FINISHED", testTask.getFinishedDTO().getEventType()); + Assertions.assertEquals("value", testTask.getFinishedDTO().getPayload().get("key")); + } + + /** + * 验证取消接口按 handleCode 路由并传递任务参数 + */ + @Test + void doHandleCancelledShouldRouteTask() { + TaskStatusCallApiReqDTO reqDTO = buildReqDTO(); + reqDTO.setEventId("1001_TASK_CANCELLED"); + reqDTO.setStatus("TASK_CANCELLED"); + + testTaskCommonApi.doHandleCancelled(reqDTO).getCheckedData(); + + Assertions.assertNotNull(testTask.getCancelledDTO()); + Assertions.assertEquals(1001L, testTask.getCancelledDTO().getTaskId()); + Assertions.assertEquals("1001_TASK_CANCELLED", testTask.getCancelledDTO().getEventId()); + Assertions.assertEquals("TASK_CANCELLED", testTask.getCancelledDTO().getEventType()); + Assertions.assertEquals("value", testTask.getCancelledDTO().getPayload().get("key")); + } + + /** + * 构建任务状态调用参数 + */ + private TaskStatusCallApiReqDTO buildReqDTO() { + TaskStatusCallApiReqDTO reqDTO = new TaskStatusCallApiReqDTO(); + reqDTO.setTaskId(1001L); + reqDTO.setTaskCode("TASK-1001"); + reqDTO.setEventId("1001_TASK_FINISHED"); + reqDTO.setStatus("TASK_FINISHED"); + reqDTO.setHandleCode("testTask"); + reqDTO.setPayload(Map.of("key", "value")); + return reqDTO; + } + + /** + * 测试任务通用接口实现 + */ + @Component + static class TestTaskCommonApi extends AbstractTaskCommonApiImpl { + } + + /** + * 测试任务处理器 + */ + @Component("testTask") + static class TestTask extends AbstractTask { + + private TaskExecuteDTO finishedDTO; + + private TaskExecuteDTO cancelledDTO; + + /** + * 记录完成任务参数 + */ + @Override + public void doHandleFinish(TaskExecuteDTO taskExecuteDTO) { + this.finishedDTO = taskExecuteDTO; + } + + /** + * 记录取消任务参数 + */ + @Override + public void doHandleCancel(TaskExecuteDTO taskExecuteDTO) { + this.cancelledDTO = taskExecuteDTO; + } + + /** + * 获取完成任务参数 + */ + public TaskExecuteDTO getFinishedDTO() { + return finishedDTO; + } + + /** + * 获取取消任务参数 + */ + public TaskExecuteDTO getCancelledDTO() { + return cancelledDTO; + } + } +} diff --git a/nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/core/TaskCommonApiFactoryLocalSpringTest.java b/nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/core/TaskCommonApiFactoryLocalSpringTest.java new file mode 100644 index 00000000..7d1ec05f --- /dev/null +++ b/nl-framework/nl-spring-boot-starter-execute/src/test/java/cn/code/nl/framework/execute/core/TaskCommonApiFactoryLocalSpringTest.java @@ -0,0 +1,61 @@ +package cn.code.nl.framework.execute.core; + +import cn.code.nl.framework.execute.biz.api.AbstractTaskCommonApiImpl; +import cn.code.nl.framework.execute.biz.api.lms.LmsTaskCommonApi; +import cn.code.nl.framework.execute.biz.api.wms.WmsTaskCommonApi; +import jakarta.annotation.Resource; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.stereotype.Component; + +/** + * 任务通用 API 工厂本地 Spring 集成测试 + */ +@SpringBootTest(classes = { + TaskFactory.class, + TaskCommonApiFactory.class, + TaskCommonApiFactoryLocalSpringTest.TestLmsTaskCommonApi.class, + TaskCommonApiFactoryLocalSpringTest.TestWmsTaskCommonApi.class +}) +class TaskCommonApiFactoryLocalSpringTest { + + @Resource + private TaskCommonApiFactory taskCommonApiFactory; + + @Resource + private TestLmsTaskCommonApi testLmsTaskCommonApi; + + @Resource + private TestWmsTaskCommonApi testWmsTaskCommonApi; + + /** + * 验证业务归属编码 LMS 能定位 Feign 接口 + */ + @Test + void getByServerNameShouldSupportLmsOwnerService() { + Assertions.assertSame(testLmsTaskCommonApi, taskCommonApiFactory.getByServerName("LMS")); + } + + /** + * 验证业务归属编码 WMS 能定位 Feign 接口 + */ + @Test + void getByServerNameShouldSupportWmsOwnerService() { + Assertions.assertSame(testWmsTaskCommonApi, taskCommonApiFactory.getByServerName("WMS")); + } + + /** + * 测试 LMS Feign 接口 + */ + @Component + static class TestLmsTaskCommonApi extends AbstractTaskCommonApiImpl implements LmsTaskCommonApi { + } + + /** + * 测试 WMS Feign 接口 + */ + @Component + static class TestWmsTaskCommonApi extends AbstractTaskCommonApiImpl implements WmsTaskCommonApi { + } +} diff --git a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/consumer/LmsTaskStatusChangeConsumer.java b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/consumer/LmsTaskStatusChangeConsumer.java index ae2959ca..0bd2a35a 100644 --- a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/consumer/LmsTaskStatusChangeConsumer.java +++ b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/consumer/LmsTaskStatusChangeConsumer.java @@ -5,17 +5,22 @@ import cn.code.nl.framework.execute.core.AbstractTask; import cn.code.nl.framework.execute.core.TaskFactory; import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO; import cn.code.nl.module.task.api.TransportTaskApi; +import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; import cn.code.nl.module.task.dto.TaskInfoDTO; +import cn.code.nl.module.task.enums.CallbackStatusEnum; import cn.code.nl.module.task.enums.TransportTaskStatusEnum; import cn.code.nl.module.task.message.TaskEventMessage; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; +import org.redisson.api.RBucket; import org.redisson.api.RLock; import org.redisson.api.RedissonClient; import org.springframework.stereotype.Component; +import java.time.Duration; + import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY; /** @@ -32,6 +37,12 @@ import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE ) public class LmsTaskStatusChangeConsumer implements RocketMQListener { + /** 任务业务处理成功标识前缀 */ + private static final String TASK_BUSINESS_HANDLED_KEY = "task:business:handled:"; + + /** 任务业务处理成功标识有效期 */ + private static final Duration TASK_BUSINESS_HANDLED_TIMEOUT = Duration.ofDays(7); + @Resource private TaskFactory taskFactory; @@ -56,7 +67,19 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener handledBucket = redissonClient.getBucket( + TASK_BUSINESS_HANDLED_KEY + message.getEventId()); + if (!Boolean.TRUE.equals(handledBucket.get())) { + try { + executeTaskHandler(message, handleCode, eventType); + handledBucket.set(true, TASK_BUSINESS_HANDLED_TIMEOUT); + } catch (RuntimeException exception) { + callbackFailedResult(message, exception); + throw exception; + } + } + callbackTaskResult(message, CallbackStatusEnum.SUCCESS.getCode(), "任务业务处理成功"); + handledBucket.delete(); } finally { if (lock.isHeldByCurrentThread()) { lock.unlock(); @@ -99,12 +122,13 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener callbackResult = transportTaskApi.receiveCallbackResult(reqDTO); + if (!callbackResult.isSuccess()) { + log.error("任务业务结果回调失败, reqDTO={}, callbackResult={}", reqDTO, callbackResult); + callbackResult.checkError(); + } + } + + /** + * 尝试回调 Task 服务记录业务处理失败 + */ + private void callbackFailedResult(TaskEventMessage message, RuntimeException exception) { + try { + callbackTaskResult(message, CallbackStatusEnum.FAILED.getCode(), exception.getMessage()); + } catch (RuntimeException callbackException) { + log.error("任务业务失败结果回调异常, message={}", message, callbackException); + } + } } diff --git a/nl-module-task/nl-module-task-api/pom.xml b/nl-module-task/nl-module-task-api/pom.xml index c6d3884c..4041a9a2 100644 --- a/nl-module-task/nl-module-task-api/pom.xml +++ b/nl-module-task/nl-module-task-api/pom.xml @@ -49,6 +49,13 @@ nl-spring-boot-starter-execute + + + org.springframework.boot + spring-boot-starter-test + test + + - \ No newline at end of file + diff --git a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskCallbackTypeEnum.java b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskCallbackTypeEnum.java new file mode 100644 index 00000000..29bd2f7a --- /dev/null +++ b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskCallbackTypeEnum.java @@ -0,0 +1,35 @@ +package cn.code.nl.module.task.enums; + +import lombok.Getter; + +/** + * 任务业务回调方式 + */ +@Getter +public enum TaskCallbackTypeEnum { + + MQ("mq"), + + FEIGN("feign"); + + private final String code; + + TaskCallbackTypeEnum(String code) { + this.code = code; + } + + /** + * 根据编码获取任务业务回调方式 + * + * @param code 回调方式编码 + * @return 任务业务回调方式 + */ + public static TaskCallbackTypeEnum getByCode(String code) { + for (TaskCallbackTypeEnum value : values()) { + if (value.getCode().equals(code)) { + return value; + } + } + return null; + } +} diff --git a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskOwnerServiceEnum.java b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskOwnerServiceEnum.java new file mode 100644 index 00000000..a9b1b386 --- /dev/null +++ b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/enums/TaskOwnerServiceEnum.java @@ -0,0 +1,50 @@ +package cn.code.nl.module.task.enums; + +import lombok.AllArgsConstructor; +import lombok.Getter; + +/** + * 任务业务归属服务枚举 + */ +@Getter +@AllArgsConstructor +public enum TaskOwnerServiceEnum { + + LMS("lms", "lms-server"), + WMS("wms", "wms-server"); + + /** + * MQ Topic 前缀 + */ + private final String topicPrefix; + + /** + * 服务注册名称 + */ + private final String serviceName; + + /** + * 根据业务归属编码获取枚举 + * + * @param code 业务归属编码 + * @return 业务归属服务枚举 + */ + public static TaskOwnerServiceEnum getByCode(String code) { + for (TaskOwnerServiceEnum value : values()) { + if (value.name().equalsIgnoreCase(code) || value.serviceName.equalsIgnoreCase(code)) { + return value; + } + } + return null; + } + + /** + * 构建业务服务消费的 Topic + * + * @param topic 基础 Topic + * @return 完整 Topic + */ + public String buildTopic(String topic) { + return topicPrefix + "_" + topic; + } +} diff --git a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/message/TaskEventMessage.java b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/message/TaskEventMessage.java index 5d7752bc..6f7efee7 100644 --- a/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/message/TaskEventMessage.java +++ b/nl-module-task/nl-module-task-api/src/main/java/cn/code/nl/module/task/message/TaskEventMessage.java @@ -44,4 +44,15 @@ public class TaskEventMessage { /** 扩展数据 */ private Map payload; + + /** + * 构建任务事件唯一标识 + * + * @param taskId 任务标识 + * @param eventType 事件类型 + * @return 任务事件唯一标识 + */ + public static String buildEventId(Long taskId, String eventType) { + return taskId + "_" + eventType; + } } diff --git a/nl-module-task/nl-module-task-api/src/test/java/cn/code/nl/module/task/enums/TaskOwnerServiceEnumLocalSpringTest.java b/nl-module-task/nl-module-task-api/src/test/java/cn/code/nl/module/task/enums/TaskOwnerServiceEnumLocalSpringTest.java new file mode 100644 index 00000000..da51f0aa --- /dev/null +++ b/nl-module-task/nl-module-task-api/src/test/java/cn/code/nl/module/task/enums/TaskOwnerServiceEnumLocalSpringTest.java @@ -0,0 +1,56 @@ +package cn.code.nl.module.task.enums; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.context.annotation.Configuration; + +/** + * 任务业务归属服务本地 Spring 集成测试 + */ +@SpringBootTest(classes = TaskOwnerServiceEnumLocalSpringTest.TestConfiguration.class) +class TaskOwnerServiceEnumLocalSpringTest { + + /** + * 验证 LMS 业务归属编码生成消费者 Topic + */ + @Test + void buildTopicShouldSupportLmsOwnerService() { + TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode("LMS"); + + Assertions.assertNotNull(ownerService); + Assertions.assertEquals("lms_task_status_change_dev_topic", + ownerService.buildTopic("task_status_change_dev_topic")); + } + + /** + * 验证 WMS 业务归属编码生成消费者 Topic + */ + @Test + void buildTopicShouldSupportWmsOwnerService() { + TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode("WMS"); + + Assertions.assertNotNull(ownerService); + Assertions.assertEquals("wms_task_status_change_dev_topic", + ownerService.buildTopic("task_status_change_dev_topic")); + } + + /** + * 验证兼容服务注册名称并拒绝非法业务归属 + */ + @Test + void getByCodeShouldSupportServiceNameAndRejectInvalidCode() { + Assertions.assertEquals(TaskOwnerServiceEnum.LMS, + TaskOwnerServiceEnum.getByCode("lms-server")); + Assertions.assertEquals(TaskOwnerServiceEnum.WMS, + TaskOwnerServiceEnum.getByCode("wms-server")); + Assertions.assertNull(TaskOwnerServiceEnum.getByCode("invalid-server")); + } + + /** + * 测试 Spring 配置 + */ + @Configuration + static class TestConfiguration { + } +} diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java index 5158daf8..4d1f3af9 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java @@ -1,22 +1,30 @@ package cn.code.nl.module.task.manage; import cn.code.nl.framework.common.exception.ServiceException; +import cn.code.nl.framework.execute.biz.api.TaskCommonApi; +import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO; +import cn.code.nl.framework.execute.core.TaskCommonApiFactory; import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO; import cn.code.nl.module.task.dal.mysql.transporttask.TransportTaskMapper; import cn.code.nl.module.task.dto.AcsFeedbackReqDTO; +import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; import cn.code.nl.module.task.enums.CallbackStatusEnum; import cn.code.nl.module.task.enums.FinishedTypeEnum; +import cn.code.nl.module.task.enums.TaskCallbackTypeEnum; import cn.code.nl.module.task.enums.TaskEventTypeEnum; import cn.code.nl.module.task.enums.TaskOperationTypeEnum; +import cn.code.nl.module.task.enums.TaskOwnerServiceEnum; import cn.code.nl.module.task.enums.TransportTaskStatusEnum; import cn.code.nl.module.task.message.TaskEventMessage; import cn.code.nl.module.task.mq.producer.TaskEventProducer; +import cn.code.nl.module.task.service.transporttask.TransportTaskCallbackService; import cn.hutool.core.util.StrUtil; import jakarta.annotation.PostConstruct; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; import org.redisson.api.RLock; import org.redisson.api.RedissonClient; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.springframework.transaction.support.TransactionSynchronization; import org.springframework.transaction.support.TransactionSynchronizationManager; @@ -51,9 +59,21 @@ public class TransportTaskOperationManager { @Resource private TransportTaskIssueManager transportTaskIssueManager; + @Resource + private TaskCommonApiFactory taskCommonApiFactory; + + @Resource + private TransportTaskCallbackService transportTaskCallbackService; + @Resource private RedissonClient redissonClient; + @Value("${nl.task.callback-type}") + private String callbackType; + + /** 任务业务回调方式 */ + private TaskCallbackTypeEnum callbackTypeEnum; + /** 操作类型路由表 */ private final Map> operationHandlers = new EnumMap<>(TaskOperationTypeEnum.class); @@ -63,6 +83,10 @@ public class TransportTaskOperationManager { */ @PostConstruct public void initOperationHandlers() { + callbackTypeEnum = TaskCallbackTypeEnum.getByCode(callbackType); + if (callbackTypeEnum == null) { + throw new IllegalArgumentException("任务业务回调方式配置错误:" + callbackType); + } operationHandlers.put(TaskOperationTypeEnum.EXECUTING, this::handleExecuting); operationHandlers.put(TaskOperationTypeEnum.ISSUE, (task, reqDTO) -> handleIssue(task)); operationHandlers.put(TaskOperationTypeEnum.FINISHED, this::handleFinished); @@ -91,7 +115,7 @@ public class TransportTaskOperationManager { throw exception(TRANSPORT_TASK_OPERATION_FAILED, ex.getMessage()); } finally { if (lock.isHeldByCurrentThread()) { - lock.unlock(); + unlockAfterTransaction(lock); } } } else { @@ -147,19 +171,7 @@ public class TransportTaskOperationManager { task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); transportTaskMapper.updateById(task); - // MQ 在事务提交后发送,避免 Consumer 读到未提交的数据 - Map payload = reqDTO.getPayload(); - if (TransactionSynchronizationManager.isSynchronizationActive()) { - TransactionSynchronizationManager.registerSynchronization( - new TransactionSynchronization() { - @Override - public void afterCommit() { - publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, payload); - } - }); - } else { - publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, payload); - } + executeBusinessCallback(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload()); } /** @@ -182,19 +194,7 @@ public class TransportTaskOperationManager { task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); transportTaskMapper.updateById(task); - // MQ 在事务提交后发送,避免 Consumer 读到未提交的数据 - Map payload = reqDTO.getPayload(); - if (TransactionSynchronizationManager.isSynchronizationActive()) { - TransactionSynchronizationManager.registerSynchronization( - new TransactionSynchronization() { - @Override - public void afterCommit() { - publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, payload); - } - }); - } else { - publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, payload); - } + executeBusinessCallback(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload()); } /** @@ -207,6 +207,104 @@ public class TransportTaskOperationManager { log.info("任务强制完成, taskId={}", task.getTaskId()); } + /** + * 在事务完成后释放任务状态锁 + */ + private void unlockAfterTransaction(RLock lock) { + if (!TransactionSynchronizationManager.isSynchronizationActive()) { + lock.unlock(); + return; + } + TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { + @Override + public void afterCompletion(int status) { + if (lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + }); + } + + /** + * 根据配置执行任务业务回调 + */ + private void executeBusinessCallback(TransportTaskDO task, TaskEventTypeEnum eventType, + Map payload) { + TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode(task.getOwnerService()); + if (ownerService == null) { + throw new ServiceException(500, "任务业务归属服务配置错误:" + task.getOwnerService()); + } + if (callbackTypeEnum == TaskCallbackTypeEnum.MQ) { + publishEventAfterCommit(task, eventType, payload); + return; + } + executeFeignCallback(task, ownerService, eventType, payload); + } + + /** + * 在事务提交后发布任务事件 + */ + private void publishEventAfterCommit(TransportTaskDO task, TaskEventTypeEnum eventType, + Map payload) { + if (TransactionSynchronizationManager.isSynchronizationActive()) { + TransactionSynchronizationManager.registerSynchronization( + new TransactionSynchronization() { + @Override + public void afterCommit() { + publishEvent(task, eventType, payload); + } + }); + return; + } + publishEvent(task, eventType, payload); + } + + /** + * 通过 Feign 同步执行任务业务并更新最终状态 + */ + private void executeFeignCallback(TransportTaskDO task, TaskOwnerServiceEnum ownerService, + TaskEventTypeEnum eventType, Map payload) { + TaskCommonApi taskCommonApi = taskCommonApiFactory.getByServerName(ownerService.name()); + TaskStatusCallApiReqDTO callApiReqDTO = buildCallApiReqDTO(task, eventType, payload); + if (eventType == TaskEventTypeEnum.TASK_FINISHED) { + taskCommonApi.doHandleFinished(callApiReqDTO).getCheckedData(); + } else { + taskCommonApi.doHandleCancelled(callApiReqDTO).getCheckedData(); + } + transportTaskCallbackService.receiveCallbackResult(buildSuccessCallbackReq(task, eventType)); + } + + /** + * 构建 Feign 任务业务调用参数 + */ + private TaskStatusCallApiReqDTO buildCallApiReqDTO(TransportTaskDO task, TaskEventTypeEnum eventType, + Map payload) { + TaskStatusCallApiReqDTO reqDTO = new TaskStatusCallApiReqDTO(); + reqDTO.setTaskId(task.getTaskId()); + reqDTO.setEventId(TaskEventMessage.buildEventId(task.getTaskId(), eventType.getCode())); + reqDTO.setTaskCode(task.getTaskCode()); + reqDTO.setStatus(eventType.getCode()); + reqDTO.setOwnerService(task.getOwnerService()); + reqDTO.setBizType(task.getBizType()); + reqDTO.setBizId(task.getBizId()); + reqDTO.setHandleCode(task.getHandleCode()); + reqDTO.setPayload(payload); + return reqDTO; + } + + /** + * 构建业务处理成功回调参数 + */ + private TaskCallbackResultReqDTO buildSuccessCallbackReq(TransportTaskDO task, + TaskEventTypeEnum eventType) { + TaskCallbackResultReqDTO reqDTO = new TaskCallbackResultReqDTO(); + reqDTO.setTaskId(task.getTaskId()); + reqDTO.setEventId(TaskEventMessage.buildEventId(task.getTaskId(), eventType.getCode())); + reqDTO.setResult(CallbackStatusEnum.SUCCESS.getCode()); + reqDTO.setMessage("任务业务处理成功"); + return reqDTO; + } + /** * 发布 MQ 事件 */ diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/producer/TaskEventProducer.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/producer/TaskEventProducer.java index 0c5c14ae..eec792c7 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/producer/TaskEventProducer.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/mq/producer/TaskEventProducer.java @@ -1,5 +1,6 @@ package cn.code.nl.module.task.mq.producer; +import cn.code.nl.module.task.enums.TaskOwnerServiceEnum; import cn.code.nl.module.task.message.TaskEventMessage; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.core.RocketMQTemplate; @@ -7,7 +8,6 @@ import org.springframework.stereotype.Component; import jakarta.annotation.Resource; import org.springframework.beans.factory.annotation.Value; -import java.util.UUID; /** * 任务事件 MQ 生产者 @@ -31,9 +31,12 @@ public class TaskEventProducer { * @param message 事件消息 */ public void publishEvent(TaskEventMessage message) { - message.setEventId(message.getTaskId().toString()); - - String destination = message.getOwnerService() + "_" + topic; + message.setEventId(TaskEventMessage.buildEventId(message.getTaskId(), message.getEventType())); + TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode(message.getOwnerService()); + if (ownerService == null) { + throw new IllegalArgumentException("任务业务归属服务配置错误:" + message.getOwnerService()); + } + String destination = ownerService.buildTopic(topic); try { rocketMQTemplate.syncSend(destination, message); log.info("MQ 发送成功, destination={}, taskId={}, eventType={}", diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java index b406f330..54758465 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/service/transporttask/TransportTaskCallbackServiceImpl.java @@ -1,24 +1,26 @@ package cn.code.nl.module.task.service.transporttask; +import cn.code.nl.framework.common.exception.ServiceException; import cn.code.nl.module.task.dal.dataobject.transporttask.TransportTaskDO; import cn.code.nl.module.task.dal.mysql.transporttask.TransportTaskMapper; import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; import cn.code.nl.module.task.enums.CallbackStatusEnum; +import cn.code.nl.module.task.enums.TaskEventTypeEnum; import cn.code.nl.module.task.enums.TransportTaskStatusEnum; +import cn.code.nl.module.task.message.TaskEventMessage; import com.mzt.logapi.context.LogRecordContext; import com.mzt.logapi.starter.annotation.LogRecord; +import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.springframework.validation.annotation.Validated; -import jakarta.annotation.Resource; - -import static cn.code.nl.module.system.enums.LogRecordConstants.*; +import static cn.code.nl.module.system.enums.LogRecordConstants.TASK_INFO; +import static cn.code.nl.module.system.enums.LogRecordConstants.TASK_INFO_CALL_BY_RPC; +import static cn.code.nl.module.system.enums.LogRecordConstants.TASK_INFO_OPERATE_TYPE; /** * 业务回调结果处理服务实现 - * - * @author 诺力管理员 */ @Slf4j @Service @@ -34,35 +36,37 @@ public class TransportTaskCallbackServiceImpl implements TransportTaskCallbackSe bizNo = "{{#reqDTO.taskId}}", success = TASK_INFO_CALL_BY_RPC) public void receiveCallbackResult(TaskCallbackResultReqDTO reqDTO) { - // 1. 查询任务 TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId()); if (task == null) { - log.warn("回调 taskId 不存在, taskId={}", reqDTO.getTaskId()); + log.warn("回调任务不存在,taskId={}", reqDTO.getTaskId()); return; } - // 2. 幂等:已是终态则直接返回 String currentStatus = task.getTaskStatus(); - String finalFinished = statusCode(TransportTaskStatusEnum.FINISHED); - String finalCancelled = statusCode(TransportTaskStatusEnum.CANCELLED); + String finalFinished = TransportTaskStatusEnum.FINISHED.getCode(); + String finalCancelled = TransportTaskStatusEnum.CANCELLED.getCode(); if (finalFinished.equals(currentStatus) || finalCancelled.equals(currentStatus)) { - log.info("任务已处终态, taskId={}, status={}, 忽略回调", + log.info("任务已处于终态,忽略重复回调,taskId={}, status={}", reqDTO.getTaskId(), currentStatus); return; } - // 3. 根据结果处理 - boolean success = "SUCCESS".equalsIgnoreCase(reqDTO.getResult()); + TaskEventTypeEnum eventType = getPendingEventType(currentStatus); + if (eventType == null) { + throw new ServiceException(500, "任务当前状态不允许处理业务回调:" + currentStatus); + } + String expectedEventId = TaskEventMessage.buildEventId(reqDTO.getTaskId(), eventType.getCode()); + if (!expectedEventId.equals(reqDTO.getEventId())) { + throw new ServiceException(500, "任务业务回调事件标识不匹配:" + reqDTO.getEventId()); + } + + boolean success = CallbackStatusEnum.SUCCESS.getCode().equalsIgnoreCase(reqDTO.getResult()); if (success) { - // 根据当前状态决定目标终态 - if (statusCode(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING).equals(currentStatus)) { - task.setTaskStatus(finalFinished); - } else if (statusCode(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING).equals(currentStatus)) { - task.setTaskStatus(finalCancelled); - } + task.setTaskStatus(eventType == TaskEventTypeEnum.TASK_FINISHED + ? finalFinished : finalCancelled); task.setCallbackStatus(CallbackStatusEnum.SUCCESS.getCode()); + task.setCallbackErrorMsg(null); } else { - // 业务处理失败:记录失败信息,增加重试次数 task.setCallbackStatus(CallbackStatusEnum.FAILED.getCode()); task.setCallbackErrorMsg(reqDTO.getMessage()); task.setCallbackRetryCount(task.getCallbackRetryCount() != null @@ -70,15 +74,22 @@ public class TransportTaskCallbackServiceImpl implements TransportTaskCallbackSe } transportTaskMapper.updateById(task); - LogRecordContext.putVariable("ownerService", task.getOwnerService()); LogRecordContext.putVariable("taskId", task.getTaskId()); - - log.info("业务回调结果已处理, taskId={}, result={}, newStatus={}", + log.info("业务回调结果已处理,taskId={}, result={}, newStatus={}", reqDTO.getTaskId(), reqDTO.getResult(), task.getTaskStatus()); } - private String statusCode(TransportTaskStatusEnum statusEnum) { - return String.format("%03d", statusEnum.getCode()); + /** + * 根据任务待处理状态获取事件类型 + */ + private TaskEventTypeEnum getPendingEventType(String currentStatus) { + if (TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING.getCode().equals(currentStatus)) { + return TaskEventTypeEnum.TASK_FINISHED; + } + if (TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING.getCode().equals(currentStatus)) { + return TaskEventTypeEnum.TASK_CANCELLED; + } + return null; } } diff --git a/nl-module-task/nl-module-task-server/src/main/resources/application-dev.yaml b/nl-module-task/nl-module-task-server/src/main/resources/application-dev.yaml index 06832135..f4e9ef62 100644 --- a/nl-module-task/nl-module-task-server/src/main/resources/application-dev.yaml +++ b/nl-module-task/nl-module-task-server/src/main/resources/application-dev.yaml @@ -7,12 +7,12 @@ spring: username: nacos password: nacos discovery: # 【注册中心】配置项 - namespace: dev # 命名空间。这里使用 dev 开发环境 + namespace: 91c8ac41-7fb0-423b-947e-2caad4c87d41 # 命名空间。这里使用 dev 开发环境 group: DEFAULT_GROUP # 使用的 Nacos 配置分组,默认为 DEFAULT_GROUP metadata: version: 1.0.0 # 服务实例的版本号,可用于灰度发布 config: # 【配置中心】配置项 - namespace: dev # 命名空间。这里使用 dev 开发环境 + namespace: 91c8ac41-7fb0-423b-947e-2caad4c87d41 # 命名空间。这里使用 dev 开发环境 group: DEFAULT_GROUP # 使用的 Nacos 配置分组,默认为 DEFAULT_GROUP --- #################### 数据库相关配置 #################### @@ -59,12 +59,12 @@ spring: master: url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例 username: root - password: root + password: root123 slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改 lazy: true # 开启懒加载,保证启动速度 url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例 username: root - password: root + password: root123 # Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优 data: @@ -83,7 +83,7 @@ rocketmq: group: task_producer_dev_group # 生产者分组 consumer: task: - topic: task-status-change-dev-topic + topic: task_status_change_dev_topic spring: # RabbitMQ 配置项,对应 RabbitProperties 配置类 @@ -100,11 +100,11 @@ spring: xxl: job: admin: - addresses: http://192.168.10.2:8080/xxl-job-admin # 调度中心部署跟地址 - accessToken: 123456 # 执行器通讯TOKEN + addresses: http://192.168.81.193:8080/xxl-job-admin # 调度中心部署跟地址 + accessToken: default_token # 执行器通讯TOKEN executor: - ip: 192.168.10.2 - port: 9994 + ip: localhost + port: 8995 --- #################### 服务保障相关配置 #################### @@ -145,4 +145,4 @@ logging: # 芋道配置项,设置当前项目所有自定义的配置 nl: - demo: false # 开启演示模式 \ No newline at end of file + demo: false # 开启演示模式 diff --git a/nl-module-task/nl-module-task-server/src/main/resources/application.yaml b/nl-module-task/nl-module-task-server/src/main/resources/application.yaml index d63db646..bfdb51d2 100644 --- a/nl-module-task/nl-module-task-server/src/main/resources/application.yaml +++ b/nl-module-task/nl-module-task-server/src/main/resources/application.yaml @@ -113,6 +113,8 @@ nl: info: version: 1.0.0 base-package: cn.code.nl.module.task + task: + callback-type: mq # 任务完成/取消业务处理方式:mq 异步、feign 同步 web: admin-ui: url: http://dashboard.nl.iocoder.cn # Admin 管理后台 UI 的地址 diff --git a/nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java b/nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java index 0c41724a..72cbba02 100644 --- a/nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java +++ b/nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java @@ -5,17 +5,22 @@ import cn.code.nl.framework.execute.core.AbstractTask; import cn.code.nl.framework.execute.core.TaskFactory; import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO; import cn.code.nl.module.task.api.TransportTaskApi; +import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; import cn.code.nl.module.task.dto.TaskInfoDTO; +import cn.code.nl.module.task.enums.CallbackStatusEnum; import cn.code.nl.module.task.enums.TransportTaskStatusEnum; import cn.code.nl.module.task.message.TaskEventMessage; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; +import org.redisson.api.RBucket; import org.redisson.api.RLock; import org.redisson.api.RedissonClient; import org.springframework.stereotype.Component; +import java.time.Duration; + import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY; /** @@ -32,6 +37,12 @@ import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE ) public class WmsTaskStatusChangeConsumer implements RocketMQListener { + /** 任务业务处理成功标识前缀 */ + private static final String TASK_BUSINESS_HANDLED_KEY = "task:business:handled:"; + + /** 任务业务处理成功标识有效期 */ + private static final Duration TASK_BUSINESS_HANDLED_TIMEOUT = Duration.ofDays(7); + @Resource private TaskFactory taskFactory; @@ -56,7 +67,19 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener handledBucket = redissonClient.getBucket( + TASK_BUSINESS_HANDLED_KEY + message.getEventId()); + if (!Boolean.TRUE.equals(handledBucket.get())) { + try { + executeTaskHandler(message, handleCode, eventType); + handledBucket.set(true, TASK_BUSINESS_HANDLED_TIMEOUT); + } catch (RuntimeException exception) { + callbackFailedResult(message, exception); + throw exception; + } + } + callbackTaskResult(message, CallbackStatusEnum.SUCCESS.getCode(), "任务业务处理成功"); + handledBucket.delete(); } finally { if (lock.isHeldByCurrentThread()) { lock.unlock(); @@ -99,12 +122,13 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener callbackResult = transportTaskApi.receiveCallbackResult(reqDTO); + if (!callbackResult.isSuccess()) { + log.error("任务业务结果回调失败, reqDTO={}, callbackResult={}", reqDTO, callbackResult); + callbackResult.checkError(); + } + } + + /** + * 尝试回调 Task 服务记录业务处理失败 + */ + private void callbackFailedResult(TaskEventMessage message, RuntimeException exception) { + try { + callbackTaskResult(message, CallbackStatusEnum.FAILED.getCode(), exception.getMessage()); + } catch (RuntimeException callbackException) { + log.error("任务业务失败结果回调异常, message={}", message, callbackException); + } + } } diff --git a/nl-server/src/main/resources/application-dev.yaml b/nl-server/src/main/resources/application-dev.yaml index 4d76d57a..9b465d8e 100644 --- a/nl-server/src/main/resources/application-dev.yaml +++ b/nl-server/src/main/resources/application-dev.yaml @@ -47,14 +47,14 @@ spring: primary: master datasource: master: - url: jdbc:mysql://192.168.10.41:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例 + url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例 username: root - password: root + password: root123 slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改 lazy: true # 开启懒加载,保证启动速度 - url: jdbc:mysql://192.168.10.41:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例 + url: jdbc:mysql://192.168.81.193:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例 username: root - password: root + password: root123 # Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优 data: diff --git a/nl-server/src/main/resources/application.yaml b/nl-server/src/main/resources/application.yaml index ec408e0f..bd810b43 100644 --- a/nl-server/src/main/resources/application.yaml +++ b/nl-server/src/main/resources/application.yaml @@ -146,13 +146,6 @@ aj: req-verify-minute-limit: 60 # verify 接口一分钟内请求数限制 --- #################### 消息队列相关 #################### - -# rocketmq 配置项,对应 RocketMQProperties 配置类 -rocketmq: - # Producer 配置项 - producer: - group: ${spring.application.name}_PRODUCER # 生产者分组 - spring: # Kafka 配置项,对应 KafkaProperties 配置类 kafka: @@ -289,6 +282,8 @@ nl: info: version: 1.0.0 base-package: cn.code.nl + task: + callback-type: mq # 任务完成/取消业务处理方式:mq 异步、feign 同步 web: admin-ui: url: http://dashboard.nl.iocoder.cn # Admin 管理后台 UI 的地址 @@ -379,4 +374,4 @@ nl: message-bus: type: redis # 消息总线的类型 -debug: false \ No newline at end of file +debug: false diff --git a/nl-server/src/test/java/cn/code/nl/server/acs/AcsUtilLocalSpringTest.java b/nl-server/src/test/java/cn/code/nl/server/acs/AcsUtilLocalSpringTest.java index f5095da7..212e1d3b 100644 --- a/nl-server/src/test/java/cn/code/nl/server/acs/AcsUtilLocalSpringTest.java +++ b/nl-server/src/test/java/cn/code/nl/server/acs/AcsUtilLocalSpringTest.java @@ -35,7 +35,7 @@ class AcsUtilLocalSpringTest { private AcsTaskDTO buildAcsTask() { AcsTaskDTO task = new AcsTaskDTO(); - task.setTaskId(10001L); +// task.setTaskId(10001L); task.setTaskCode("LOCAL-ACS-TEST-10001"); task.setStartDeviceCode("START-TEST-01"); task.setNextDeviceCode("END-TEST-01");