Merge branch 'refs/heads/feature/20260713/task-module' into feature/20260808/packaging_workshop

# Conflicts:
#	nl-framework/nl-spring-boot-starter-execute/src/main/java/cn/code/nl/framework/execute/util/http/AcsUtil.java
#	nl-server/src/test/java/cn/code/nl/server/acs/AcsUtilLocalSpringTest.java
This commit is contained in:
2026-08-10 14:28:27 +08:00
25 changed files with 834 additions and 81 deletions

View File

@@ -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` 通过,仅有工作区行尾转换提示。

View File

@@ -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 本地事务不构成分布式事务;远端成功后若本地事务回滚,重试仍依赖上述事件幂等。

View File

@@ -47,6 +47,13 @@
<artifactId>springdoc-openapi-starter-webmvc-ui</artifactId> <artifactId>springdoc-openapi-starter-webmvc-ui</artifactId>
<scope>provided</scope> <!-- 设置为 provided主要是 PageParam 使用到 --> <scope>provided</scope> <!-- 设置为 provided主要是 PageParam 使用到 -->
</dependency> </dependency>
<!-- Test 测试相关 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies> </dependencies>
</project> </project>

View File

@@ -22,6 +22,29 @@ public abstract class AbstractTaskCommonApiImpl implements TaskCommonApi {
@Resource @Resource
private TaskFactory taskFactory; private TaskFactory taskFactory;
/**
* 处理任务完成
*/
@Override
public CommonResult<AcsApplyActionRespVO> doHandleFinished(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) {
return execute(taskStatusCallApiReqDTO, (task, taskExecuteDTO) -> {
task.doHandleFinish(taskExecuteDTO);
return null;
});
}
/**
* 处理任务取消
*/
@Override
public CommonResult<AcsApplyActionRespVO> 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) { private TaskExecuteDTO buildTaskExecuteDTO(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) {
return TaskExecuteDTO.builder() return TaskExecuteDTO.builder()
.taskId(taskStatusCallApiReqDTO.getTaskId()) .taskId(taskStatusCallApiReqDTO.getTaskId())
.eventId(taskStatusCallApiReqDTO.getEventId())
.eventType(taskStatusCallApiReqDTO.getStatus())
.payload(taskStatusCallApiReqDTO.getPayload()) .payload(taskStatusCallApiReqDTO.getPayload())
.build(); .build();
} }

View File

@@ -14,6 +14,14 @@ import org.springframework.web.bind.annotation.RequestBody;
*/ */
public interface TaskCommonApi { public interface TaskCommonApi {
@Operation(summary = "处理任务完成")
@PostMapping("/doHandleFinished")
CommonResult<AcsApplyActionRespVO> doHandleFinished(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO);
@Operation(summary = "处理任务取消")
@PostMapping("/doHandleCancelled")
CommonResult<AcsApplyActionRespVO> doHandleCancelled(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO);
@Operation(summary = "请求取货") @Operation(summary = "请求取货")
@PostMapping("/do-handle-picked") @PostMapping("/do-handle-picked")
CommonResult<AcsApplyActionRespVO> doHandlePicked(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO); CommonResult<AcsApplyActionRespVO> doHandlePicked(@RequestBody TaskStatusCallApiReqDTO taskStatusCallApiReqDTO);

View File

@@ -15,6 +15,9 @@ public class TaskStatusCallApiReqDTO implements Serializable {
/** 任务ID */ /** 任务ID */
private Long taskId; private Long taskId;
/** 任务事件唯一标识 */
private String eventId;
/** 任务编码 */ /** 任务编码 */
private String taskCode; private String taskCode;

View File

@@ -23,6 +23,12 @@ import java.util.Map;
@Component @Component
public class TaskCommonApiFactory { public class TaskCommonApiFactory {
/** LMS 业务归属编码 */
private static final String LMS_OWNER_SERVICE = "LMS";
/** WMS 业务归属编码 */
private static final String WMS_OWNER_SERVICE = "WMS";
@Autowired(required = false) @Autowired(required = false)
private LmsTaskCommonApi lmsTaskCommonApi; private LmsTaskCommonApi lmsTaskCommonApi;
@@ -41,9 +47,11 @@ public class TaskCommonApiFactory {
public void init() { public void init() {
if (lmsTaskCommonApi != null) { if (lmsTaskCommonApi != null) {
apiMap.put(RpcConstants.LMS_NAME, lmsTaskCommonApi); apiMap.put(RpcConstants.LMS_NAME, lmsTaskCommonApi);
apiMap.put(LMS_OWNER_SERVICE, lmsTaskCommonApi);
} }
if (wmsTaskCommonApi != null) { if (wmsTaskCommonApi != null) {
apiMap.put(RpcConstants.WMS_NAME, wmsTaskCommonApi); apiMap.put(RpcConstants.WMS_NAME, wmsTaskCommonApi);
apiMap.put(WMS_OWNER_SERVICE, wmsTaskCommonApi);
} }
} }

View File

@@ -24,6 +24,12 @@ public class TaskExecuteDTO {
@NotNull(message = "任务ID不能为空") @NotNull(message = "任务ID不能为空")
private Long taskId; private Long taskId;
@Schema(description = "任务事件唯一标识")
private String eventId;
@Schema(description = "任务事件类型")
private String eventType;
@Schema(description = "ACS反馈扩展数据") @Schema(description = "ACS反馈扩展数据")
private Map<String, Object> payload; private Map<String, Object> payload;
} }

View File

@@ -53,7 +53,7 @@ public class AcsUtil {
String response; String response;
try { try {
response = HttpUtils.post(url, headers(), JsonUtils.toJsonString(request)); response = HttpUtils.post(url, headers(), JsonUtils.toJsonString(request));
log.info("ACS 返回值:{}", JSON.toJSONString(response)); log.info("ACS 返回值:{}", response);
} catch (RuntimeException ex) { } catch (RuntimeException ex) {
if (isNetworkException(ex)) { if (isNetworkException(ex)) {
log.error("ACS 服务网络不通url={}request={}", url, JsonUtils.toJsonString(request), ex); log.error("ACS 服务网络不通url={}request={}", url, JsonUtils.toJsonString(request), ex);

View File

@@ -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;
}
}
}

View File

@@ -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 {
}
}

View File

@@ -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.TaskFactory;
import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO; import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO;
import cn.code.nl.module.task.api.TransportTaskApi; 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.dto.TaskInfoDTO;
import cn.code.nl.module.task.enums.CallbackStatusEnum;
import cn.code.nl.module.task.enums.TransportTaskStatusEnum; import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
import cn.code.nl.module.task.message.TaskEventMessage; import cn.code.nl.module.task.message.TaskEventMessage;
import jakarta.annotation.Resource; import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener; import org.apache.rocketmq.spring.core.RocketMQListener;
import org.redisson.api.RBucket;
import org.redisson.api.RLock; import org.redisson.api.RLock;
import org.redisson.api.RedissonClient; import org.redisson.api.RedissonClient;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import java.time.Duration;
import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY; 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<TaskEventMessage> { public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMessage> {
/** 任务业务处理成功标识前缀 */
private static final String TASK_BUSINESS_HANDLED_KEY = "task:business:handled:";
/** 任务业务处理成功标识有效期 */
private static final Duration TASK_BUSINESS_HANDLED_TIMEOUT = Duration.ofDays(7);
@Resource @Resource
private TaskFactory taskFactory; private TaskFactory taskFactory;
@@ -56,7 +67,19 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
if (!isTaskCallbackPending(message)) { if (!isTaskCallbackPending(message)) {
return; return;
} }
executeTaskHandler(message, handleCode, eventType); RBucket<Boolean> 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 { } finally {
if (lock.isHeldByCurrentThread()) { if (lock.isHeldByCurrentThread()) {
lock.unlock(); lock.unlock();
@@ -99,12 +122,13 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) { private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) {
AbstractTask task = taskFactory.getTask(handleCode); AbstractTask task = taskFactory.getTask(handleCode);
if (task == null) { if (task == null) {
log.warn("未找到对应任务处理器, handleCode={}", handleCode); throw new IllegalStateException("未找到对应任务处理器:" + handleCode);
return;
} }
TaskExecuteDTO dto = new TaskExecuteDTO(); TaskExecuteDTO dto = new TaskExecuteDTO();
dto.setTaskId(message.getTaskId()); dto.setTaskId(message.getTaskId());
dto.setEventId(message.getEventId());
dto.setEventType(message.getEventType());
dto.setPayload(message.getPayload()); dto.setPayload(message.getPayload());
if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) { if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) {
@@ -115,4 +139,31 @@ public class LmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
log.warn("未知事件类型, eventType={}, handleCode={}", eventType, handleCode); log.warn("未知事件类型, eventType={}, handleCode={}", eventType, handleCode);
} }
} }
/**
* 回调 Task 服务记录业务处理结果
*/
private void callbackTaskResult(TaskEventMessage message, String result, String resultMessage) {
TaskCallbackResultReqDTO reqDTO = new TaskCallbackResultReqDTO();
reqDTO.setTaskId(message.getTaskId());
reqDTO.setEventId(message.getEventId());
reqDTO.setResult(result);
reqDTO.setMessage(resultMessage);
CommonResult<Boolean> 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);
}
}
} }

View File

@@ -49,6 +49,13 @@
<artifactId>nl-spring-boot-starter-execute</artifactId> <artifactId>nl-spring-boot-starter-execute</artifactId>
</dependency> </dependency>
<!-- Test 测试相关 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies> </dependencies>
</project> </project>

View File

@@ -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;
}
}

View File

@@ -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;
}
}

View File

@@ -44,4 +44,15 @@ public class TaskEventMessage {
/** 扩展数据 */ /** 扩展数据 */
private Map<String, Object> payload; private Map<String, Object> payload;
/**
* 构建任务事件唯一标识
*
* @param taskId 任务标识
* @param eventType 事件类型
* @return 任务事件唯一标识
*/
public static String buildEventId(Long taskId, String eventType) {
return taskId + "_" + eventType;
}
} }

View File

@@ -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 {
}
}

View File

@@ -1,22 +1,30 @@
package cn.code.nl.module.task.manage; package cn.code.nl.module.task.manage;
import cn.code.nl.framework.common.exception.ServiceException; 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.dataobject.transporttask.TransportTaskDO;
import cn.code.nl.module.task.dal.mysql.transporttask.TransportTaskMapper; 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.AcsFeedbackReqDTO;
import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO;
import cn.code.nl.module.task.enums.CallbackStatusEnum; import cn.code.nl.module.task.enums.CallbackStatusEnum;
import cn.code.nl.module.task.enums.FinishedTypeEnum; 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.TaskEventTypeEnum;
import cn.code.nl.module.task.enums.TaskOperationTypeEnum; 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.enums.TransportTaskStatusEnum;
import cn.code.nl.module.task.message.TaskEventMessage; import cn.code.nl.module.task.message.TaskEventMessage;
import cn.code.nl.module.task.mq.producer.TaskEventProducer; 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 cn.hutool.core.util.StrUtil;
import jakarta.annotation.PostConstruct; import jakarta.annotation.PostConstruct;
import jakarta.annotation.Resource; import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.redisson.api.RLock; import org.redisson.api.RLock;
import org.redisson.api.RedissonClient; import org.redisson.api.RedissonClient;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.springframework.transaction.support.TransactionSynchronization; import org.springframework.transaction.support.TransactionSynchronization;
import org.springframework.transaction.support.TransactionSynchronizationManager; import org.springframework.transaction.support.TransactionSynchronizationManager;
@@ -51,9 +59,21 @@ public class TransportTaskOperationManager {
@Resource @Resource
private TransportTaskIssueManager transportTaskIssueManager; private TransportTaskIssueManager transportTaskIssueManager;
@Resource
private TaskCommonApiFactory taskCommonApiFactory;
@Resource
private TransportTaskCallbackService transportTaskCallbackService;
@Resource @Resource
private RedissonClient redissonClient; private RedissonClient redissonClient;
@Value("${nl.task.callback-type}")
private String callbackType;
/** 任务业务回调方式 */
private TaskCallbackTypeEnum callbackTypeEnum;
/** 操作类型路由表 */ /** 操作类型路由表 */
private final Map<TaskOperationTypeEnum, BiConsumer<TransportTaskDO, AcsFeedbackReqDTO>> operationHandlers = private final Map<TaskOperationTypeEnum, BiConsumer<TransportTaskDO, AcsFeedbackReqDTO>> operationHandlers =
new EnumMap<>(TaskOperationTypeEnum.class); new EnumMap<>(TaskOperationTypeEnum.class);
@@ -63,6 +83,10 @@ public class TransportTaskOperationManager {
*/ */
@PostConstruct @PostConstruct
public void initOperationHandlers() { public void initOperationHandlers() {
callbackTypeEnum = TaskCallbackTypeEnum.getByCode(callbackType);
if (callbackTypeEnum == null) {
throw new IllegalArgumentException("任务业务回调方式配置错误:" + callbackType);
}
operationHandlers.put(TaskOperationTypeEnum.EXECUTING, this::handleExecuting); operationHandlers.put(TaskOperationTypeEnum.EXECUTING, this::handleExecuting);
operationHandlers.put(TaskOperationTypeEnum.ISSUE, (task, reqDTO) -> handleIssue(task)); operationHandlers.put(TaskOperationTypeEnum.ISSUE, (task, reqDTO) -> handleIssue(task));
operationHandlers.put(TaskOperationTypeEnum.FINISHED, this::handleFinished); operationHandlers.put(TaskOperationTypeEnum.FINISHED, this::handleFinished);
@@ -91,7 +115,7 @@ public class TransportTaskOperationManager {
throw exception(TRANSPORT_TASK_OPERATION_FAILED, ex.getMessage()); throw exception(TRANSPORT_TASK_OPERATION_FAILED, ex.getMessage());
} finally { } finally {
if (lock.isHeldByCurrentThread()) { if (lock.isHeldByCurrentThread()) {
lock.unlock(); unlockAfterTransaction(lock);
} }
} }
} else { } else {
@@ -147,19 +171,7 @@ public class TransportTaskOperationManager {
task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode());
transportTaskMapper.updateById(task); transportTaskMapper.updateById(task);
// MQ 在事务提交后发送,避免 Consumer 读到未提交的数据 executeBusinessCallback(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload());
Map<String, Object> 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);
}
} }
/** /**
@@ -182,19 +194,7 @@ public class TransportTaskOperationManager {
task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode());
transportTaskMapper.updateById(task); transportTaskMapper.updateById(task);
// MQ 在事务提交后发送,避免 Consumer 读到未提交的数据 executeBusinessCallback(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload());
Map<String, Object> 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);
}
} }
/** /**
@@ -207,6 +207,104 @@ public class TransportTaskOperationManager {
log.info("任务强制完成, taskId={}", task.getTaskId()); 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<String, Object> 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<String, Object> 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<String, Object> 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<String, Object> 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 事件 * 发布 MQ 事件
*/ */

View File

@@ -1,5 +1,6 @@
package cn.code.nl.module.task.mq.producer; 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 cn.code.nl.module.task.message.TaskEventMessage;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.apache.rocketmq.spring.core.RocketMQTemplate;
@@ -7,7 +8,6 @@ import org.springframework.stereotype.Component;
import jakarta.annotation.Resource; import jakarta.annotation.Resource;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import java.util.UUID;
/** /**
* 任务事件 MQ 生产者 * 任务事件 MQ 生产者
@@ -31,9 +31,12 @@ public class TaskEventProducer {
* @param message 事件消息 * @param message 事件消息
*/ */
public void publishEvent(TaskEventMessage message) { public void publishEvent(TaskEventMessage message) {
message.setEventId(message.getTaskId().toString()); message.setEventId(TaskEventMessage.buildEventId(message.getTaskId(), message.getEventType()));
TaskOwnerServiceEnum ownerService = TaskOwnerServiceEnum.getByCode(message.getOwnerService());
String destination = message.getOwnerService() + "_" + topic; if (ownerService == null) {
throw new IllegalArgumentException("任务业务归属服务配置错误:" + message.getOwnerService());
}
String destination = ownerService.buildTopic(topic);
try { try {
rocketMQTemplate.syncSend(destination, message); rocketMQTemplate.syncSend(destination, message);
log.info("MQ 发送成功, destination={}, taskId={}, eventType={}", log.info("MQ 发送成功, destination={}, taskId={}, eventType={}",

View File

@@ -1,24 +1,26 @@
package cn.code.nl.module.task.service.transporttask; 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.dataobject.transporttask.TransportTaskDO;
import cn.code.nl.module.task.dal.mysql.transporttask.TransportTaskMapper; import cn.code.nl.module.task.dal.mysql.transporttask.TransportTaskMapper;
import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO; import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO;
import cn.code.nl.module.task.enums.CallbackStatusEnum; 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.enums.TransportTaskStatusEnum;
import cn.code.nl.module.task.message.TaskEventMessage;
import com.mzt.logapi.context.LogRecordContext; import com.mzt.logapi.context.LogRecordContext;
import com.mzt.logapi.starter.annotation.LogRecord; import com.mzt.logapi.starter.annotation.LogRecord;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.validation.annotation.Validated; import org.springframework.validation.annotation.Validated;
import jakarta.annotation.Resource; 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.*; import static cn.code.nl.module.system.enums.LogRecordConstants.TASK_INFO_OPERATE_TYPE;
/** /**
* 业务回调结果处理服务实现 * 业务回调结果处理服务实现
*
* @author 诺力管理员
*/ */
@Slf4j @Slf4j
@Service @Service
@@ -34,35 +36,37 @@ public class TransportTaskCallbackServiceImpl implements TransportTaskCallbackSe
bizNo = "{{#reqDTO.taskId}}", bizNo = "{{#reqDTO.taskId}}",
success = TASK_INFO_CALL_BY_RPC) success = TASK_INFO_CALL_BY_RPC)
public void receiveCallbackResult(TaskCallbackResultReqDTO reqDTO) { public void receiveCallbackResult(TaskCallbackResultReqDTO reqDTO) {
// 1. 查询任务
TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId()); TransportTaskDO task = transportTaskMapper.selectById(reqDTO.getTaskId());
if (task == null) { if (task == null) {
log.warn("回调 taskId 不存在, taskId={}", reqDTO.getTaskId()); log.warn("回调任务不存在taskId={}", reqDTO.getTaskId());
return; return;
} }
// 2. 幂等:已是终态则直接返回
String currentStatus = task.getTaskStatus(); String currentStatus = task.getTaskStatus();
String finalFinished = statusCode(TransportTaskStatusEnum.FINISHED); String finalFinished = TransportTaskStatusEnum.FINISHED.getCode();
String finalCancelled = statusCode(TransportTaskStatusEnum.CANCELLED); String finalCancelled = TransportTaskStatusEnum.CANCELLED.getCode();
if (finalFinished.equals(currentStatus) || finalCancelled.equals(currentStatus)) { if (finalFinished.equals(currentStatus) || finalCancelled.equals(currentStatus)) {
log.info("任务已处终态, taskId={}, status={}, 忽略回调", log.info("任务已处终态,忽略重复回调,taskId={}, status={}",
reqDTO.getTaskId(), currentStatus); reqDTO.getTaskId(), currentStatus);
return; return;
} }
// 3. 根据结果处理 TaskEventTypeEnum eventType = getPendingEventType(currentStatus);
boolean success = "SUCCESS".equalsIgnoreCase(reqDTO.getResult()); 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 (success) {
// 根据当前状态决定目标终态 task.setTaskStatus(eventType == TaskEventTypeEnum.TASK_FINISHED
if (statusCode(TransportTaskStatusEnum.FINISHED_CALLBACK_PENDING).equals(currentStatus)) { ? finalFinished : finalCancelled);
task.setTaskStatus(finalFinished);
} else if (statusCode(TransportTaskStatusEnum.CANCEL_CALLBACK_PENDING).equals(currentStatus)) {
task.setTaskStatus(finalCancelled);
}
task.setCallbackStatus(CallbackStatusEnum.SUCCESS.getCode()); task.setCallbackStatus(CallbackStatusEnum.SUCCESS.getCode());
task.setCallbackErrorMsg(null);
} else { } else {
// 业务处理失败:记录失败信息,增加重试次数
task.setCallbackStatus(CallbackStatusEnum.FAILED.getCode()); task.setCallbackStatus(CallbackStatusEnum.FAILED.getCode());
task.setCallbackErrorMsg(reqDTO.getMessage()); task.setCallbackErrorMsg(reqDTO.getMessage());
task.setCallbackRetryCount(task.getCallbackRetryCount() != null task.setCallbackRetryCount(task.getCallbackRetryCount() != null
@@ -70,15 +74,22 @@ public class TransportTaskCallbackServiceImpl implements TransportTaskCallbackSe
} }
transportTaskMapper.updateById(task); transportTaskMapper.updateById(task);
LogRecordContext.putVariable("ownerService", task.getOwnerService()); LogRecordContext.putVariable("ownerService", task.getOwnerService());
LogRecordContext.putVariable("taskId", task.getTaskId()); LogRecordContext.putVariable("taskId", task.getTaskId());
log.info("业务回调结果已处理taskId={}, result={}, newStatus={}",
log.info("业务回调结果已处理, taskId={}, result={}, newStatus={}",
reqDTO.getTaskId(), reqDTO.getResult(), task.getTaskStatus()); 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;
} }
} }

View File

@@ -59,12 +59,12 @@ spring:
master: master:
url: jdbc:mysql://127.0.0.1:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例 url: jdbc:mysql://127.0.0.1:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
username: root username: root
password: root password: root123
slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改 slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改
lazy: true # 开启懒加载,保证启动速度 lazy: true # 开启懒加载,保证启动速度
url: jdbc:mysql://127.0.0.1:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例 url: jdbc:mysql://127.0.0.1:3306/huachuang_lms_dev?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&nullCatalogMeansCurrent=true&rewriteBatchedStatements=true # MySQL Connector/J 8.X 连接的示例
username: root username: root
password: root password: root123
# Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优 # Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优
data: data:
@@ -83,7 +83,7 @@ rocketmq:
group: task_producer_dev_group # 生产者分组 group: task_producer_dev_group # 生产者分组
consumer: consumer:
task: task:
topic: task-status-change-dev-topic topic: task_status_change_dev_topic
spring: spring:
# RabbitMQ 配置项,对应 RabbitProperties 配置类 # RabbitMQ 配置项,对应 RabbitProperties 配置类
@@ -145,4 +145,4 @@ logging:
# 芋道配置项,设置当前项目所有自定义的配置 # 芋道配置项,设置当前项目所有自定义的配置
nl: nl:
demo: false # 开启演示模式 demo: false # 开启演示模式

View File

@@ -113,6 +113,8 @@ nl:
info: info:
version: 1.0.0 version: 1.0.0
base-package: cn.code.nl.module.task base-package: cn.code.nl.module.task
task:
callback-type: mq # 任务完成/取消业务处理方式mq 异步、feign 同步
web: web:
admin-ui: admin-ui:
url: http://dashboard.nl.iocoder.cn # Admin 管理后台 UI 的地址 url: http://dashboard.nl.iocoder.cn # Admin 管理后台 UI 的地址

View File

@@ -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.TaskFactory;
import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO; import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO;
import cn.code.nl.module.task.api.TransportTaskApi; 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.dto.TaskInfoDTO;
import cn.code.nl.module.task.enums.CallbackStatusEnum;
import cn.code.nl.module.task.enums.TransportTaskStatusEnum; import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
import cn.code.nl.module.task.message.TaskEventMessage; import cn.code.nl.module.task.message.TaskEventMessage;
import jakarta.annotation.Resource; import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener; import org.apache.rocketmq.spring.core.RocketMQListener;
import org.redisson.api.RBucket;
import org.redisson.api.RLock; import org.redisson.api.RLock;
import org.redisson.api.RedissonClient; import org.redisson.api.RedissonClient;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import java.time.Duration;
import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY; 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<TaskEventMessage> { public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMessage> {
/** 任务业务处理成功标识前缀 */
private static final String TASK_BUSINESS_HANDLED_KEY = "task:business:handled:";
/** 任务业务处理成功标识有效期 */
private static final Duration TASK_BUSINESS_HANDLED_TIMEOUT = Duration.ofDays(7);
@Resource @Resource
private TaskFactory taskFactory; private TaskFactory taskFactory;
@@ -56,7 +67,19 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
if (!isTaskCallbackPending(message)) { if (!isTaskCallbackPending(message)) {
return; return;
} }
executeTaskHandler(message, handleCode, eventType); RBucket<Boolean> 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 { } finally {
if (lock.isHeldByCurrentThread()) { if (lock.isHeldByCurrentThread()) {
lock.unlock(); lock.unlock();
@@ -99,12 +122,13 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) { private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) {
AbstractTask task = taskFactory.getTask(handleCode); AbstractTask task = taskFactory.getTask(handleCode);
if (task == null) { if (task == null) {
log.warn("未找到对应任务处理器, handleCode={}", handleCode); throw new IllegalStateException("未找到对应任务处理器:" + handleCode);
return;
} }
TaskExecuteDTO dto = new TaskExecuteDTO(); TaskExecuteDTO dto = new TaskExecuteDTO();
dto.setTaskId(message.getTaskId()); dto.setTaskId(message.getTaskId());
dto.setEventId(message.getEventId());
dto.setEventType(message.getEventType());
dto.setPayload(message.getPayload()); dto.setPayload(message.getPayload());
if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) { if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) {
@@ -115,4 +139,31 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
log.warn("未知事件类型, eventType={}, handleCode={}", eventType, handleCode); log.warn("未知事件类型, eventType={}, handleCode={}", eventType, handleCode);
} }
} }
/**
* 回调 Task 服务记录业务处理结果
*/
private void callbackTaskResult(TaskEventMessage message, String result, String resultMessage) {
TaskCallbackResultReqDTO reqDTO = new TaskCallbackResultReqDTO();
reqDTO.setTaskId(message.getTaskId());
reqDTO.setEventId(message.getEventId());
reqDTO.setResult(result);
reqDTO.setMessage(resultMessage);
CommonResult<Boolean> 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);
}
}
} }

View File

@@ -47,14 +47,14 @@ spring:
primary: master primary: master
datasource: datasource:
master: 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 username: root
password: root password: root123
slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改 slave: # 模拟从库,可根据自己需要修改 # 模拟从库,可根据自己需要修改
lazy: true # 开启懒加载,保证启动速度 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 username: root
password: root password: root123
# Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优 # Redis 配置。Redisson 默认的配置足够使用,一般不需要进行调优
data: data:

View File

@@ -146,13 +146,6 @@ aj:
req-verify-minute-limit: 60 # verify 接口一分钟内请求数限制 req-verify-minute-limit: 60 # verify 接口一分钟内请求数限制
--- #################### 消息队列相关 #################### --- #################### 消息队列相关 ####################
# rocketmq 配置项,对应 RocketMQProperties 配置类
rocketmq:
# Producer 配置项
producer:
group: ${spring.application.name}_PRODUCER # 生产者分组
spring: spring:
# Kafka 配置项,对应 KafkaProperties 配置类 # Kafka 配置项,对应 KafkaProperties 配置类
kafka: kafka:
@@ -289,6 +282,8 @@ nl:
info: info:
version: 1.0.0 version: 1.0.0
base-package: cn.code.nl base-package: cn.code.nl
task:
callback-type: mq # 任务完成/取消业务处理方式mq 异步、feign 同步
web: web:
admin-ui: admin-ui:
url: http://dashboard.nl.iocoder.cn # Admin 管理后台 UI 的地址 url: http://dashboard.nl.iocoder.cn # Admin 管理后台 UI 的地址
@@ -379,4 +374,4 @@ nl:
message-bus: message-bus:
type: redis # 消息总线的类型 type: redis # 消息总线的类型
debug: false debug: false