From 6f566b7a694e7ff7d3efa1780947e710bfa611b7 Mon Sep 17 00:00:00 2001 From: liyongde <1419499670@qq.com> Date: Thu, 23 Jul 2026 16:48:05 +0800 Subject: [PATCH] =?UTF-8?q?opt:=20=E4=BB=BB=E5=8A=A1=E5=8F=91=E9=80=81?= =?UTF-8?q?=E6=B6=88=E8=B4=B9=E8=80=85topic=E5=AE=9A=E4=B9=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- nl-module-lms/nl-module-lms-server/pom.xml | 6 +- .../cn/code/nl/module/lms/demo/DemoApi.java | 80 ------------------- .../consumer/LmsTaskStatusChangeConsumer.java | 54 +++++++++++++ .../code/nl/module/lms/mq/package-info.java | 6 ++ .../main/resources/application-mq-dev.yaml | 11 +++ .../src/main/resources/application.yaml | 1 + .../task/mq/producer/TaskEventProducer.java | 8 +- .../src/main/resources/application-dev.yaml | 7 +- ....java => WmsTaskStatusChangeConsumer.java} | 2 +- .../main/resources/application-mq-dev.yaml | 3 + .../main/resources/application-mq-dev.yaml | 3 + 11 files changed, 94 insertions(+), 87 deletions(-) delete mode 100644 nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoApi.java create mode 100644 nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/consumer/LmsTaskStatusChangeConsumer.java create mode 100644 nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/package-info.java create mode 100644 nl-module-lms/nl-module-lms-server/src/main/resources/application-mq-dev.yaml rename nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/{TaskStatusChangeConsumer.java => WmsTaskStatusChangeConsumer.java} (95%) diff --git a/nl-module-lms/nl-module-lms-server/pom.xml b/nl-module-lms/nl-module-lms-server/pom.xml index 562d3f5f..b415622a 100644 --- a/nl-module-lms/nl-module-lms-server/pom.xml +++ b/nl-module-lms/nl-module-lms-server/pom.xml @@ -32,7 +32,11 @@ nl-module-lms-api ${revision} - + + cn.nl.cloud + nl-module-task-api + ${revision} + cn.nl.cloud diff --git a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoApi.java b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoApi.java deleted file mode 100644 index b08f3ad9..00000000 --- a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/demo/DemoApi.java +++ /dev/null @@ -1,80 +0,0 @@ -//package cn.code.nl.module.lms.demo; -// -//import cn.code.nl.framework.common.pojo.CommonResult; -//import cn.code.nl.framework.execute.biz.api.lms.LmsTaskCommonApi; -//import cn.code.nl.framework.execute.biz.dto.TaskStatusCallApiReqDTO; -//import cn.code.nl.framework.execute.biz.vo.AcsApplyActionRespVO; -//import com.alibaba.fastjson.JSON; -//import org.springframework.validation.annotation.Validated; -//import org.springframework.web.bind.annotation.RestController; -// -//import static cn.code.nl.framework.common.pojo.CommonResult.success; -// -///** -// * -// * @Author: liyongde -// * @Date: 2026/7/15 15:25 -// */ -//@RestController // 提供 RESTful API 接口,给 Feign 调用 -//@Validated -//public class DemoApi implements LmsTaskCommonApi { -// @Override -// public CommonResult doHandlePicked(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { -// // 构建外层VO -// AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); -// respVO.setTaskId(10001L); -// respVO.setTaskCode("ACS20260716001"); -// // 内层data赋值 "success" -// return success(respVO); -// } -// -// @Override -// public CommonResult doHandleApplyAgain(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { -// DemoRespVO demoRespVO = new DemoRespVO(); -// demoRespVO.setTargetPoint("A_10001"); -// -// // 构建外层VO -// AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); -// respVO.setTaskId(10001L); -// respVO.setTaskCode("ACS20260716001"); -// // 内层data赋值 "success" -// respVO.setData(JSON.parseObject(JSON.toJSONString(demoRespVO))); -// return success(respVO); -// } -// -// @Override -// public CommonResult doHandleRequestRelease(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { -// // 构建外层VO -// AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); -// respVO.setTaskId(10001L); -// respVO.setTaskCode("ACS20260716001"); -// return success(respVO); -// } -// -// @Override -// public CommonResult doHandleRequestPick(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { -// // 构建外层VO -// AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); -// respVO.setTaskId(10001L); -// respVO.setTaskCode("ACS20260716001"); -// return success(respVO); -// } -// -// @Override -// public CommonResult doHandleRequestLeave(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { -// // 构建外层VO -// AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); -// respVO.setTaskId(10001L); -// respVO.setTaskCode("ACS20260716001"); -// return success(respVO); -// } -// -// @Override -// public CommonResult doHandleRequestEnter(TaskStatusCallApiReqDTO taskStatusCallApiReqDTO) { -// // 构建外层VO -// AcsApplyActionRespVO respVO = new AcsApplyActionRespVO(); -// respVO.setTaskId(10001L); -// respVO.setTaskCode("ACS20260716001"); -// return success(respVO); -// } -//} 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 new file mode 100644 index 00000000..58a38267 --- /dev/null +++ b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/consumer/LmsTaskStatusChangeConsumer.java @@ -0,0 +1,54 @@ +package cn.code.nl.module.lms.mq.consumer; + +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.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.springframework.stereotype.Component; + +/** + * 监听任务状态变更,根据 handleCode 定位任务子类并执行完成/取消逻辑 + * + * @Author: liyongde + * @Date: 2026/7/20 10:28 + */ +@Slf4j +@Component +@RocketMQMessageListener( + topic = "${rocketmq.consumer.lms-task-operate.topic}", + consumerGroup = "${rocketmq.consumer.lms-task-operate.group}" +) +public class LmsTaskStatusChangeConsumer implements RocketMQListener { + + @Resource + private TaskFactory taskFactory; + + @Override + public void onMessage(TaskEventMessage message) { + String handleCode = message.getHandleCode(); + String eventType = message.getEventType(); + log.info("收到任务状态变更消息, handleCode={}, eventType={}, taskId={}", handleCode, eventType, message.getTaskId()); + + AbstractTask task = taskFactory.getTask(handleCode); + if (task == null) { + log.warn("未找到对应任务处理器, handleCode={}", handleCode); + return; + } + + TaskExecuteDTO dto = new TaskExecuteDTO(); + dto.setTaskId(message.getTaskId()); + dto.setPayload(message.getPayload()); + + if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) { + task.doHandleFinish(dto); + } else if (TaskEventMessage.EVENT_TYPE_CANCELLED.equals(eventType)) { + task.doHandleCancel(dto); + } else { + log.warn("未知事件类型, eventType={}, handleCode={}", eventType, handleCode); + } + } +} diff --git a/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/package-info.java b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/package-info.java new file mode 100644 index 00000000..dfa0616a --- /dev/null +++ b/nl-module-lms/nl-module-lms-server/src/main/java/cn/code/nl/module/lms/mq/package-info.java @@ -0,0 +1,6 @@ +/** + * + * @Author: liyongde + * @Date: 2026/7/23 16:21 + */ +package cn.code.nl.module.lms.mq; \ No newline at end of file diff --git a/nl-module-lms/nl-module-lms-server/src/main/resources/application-mq-dev.yaml b/nl-module-lms/nl-module-lms-server/src/main/resources/application-mq-dev.yaml new file mode 100644 index 00000000..3fab2cc8 --- /dev/null +++ b/nl-module-lms/nl-module-lms-server/src/main/resources/application-mq-dev.yaml @@ -0,0 +1,11 @@ +--- #################### MQ 消息队列相关配置 #################### +# rocketmq 配置项,对应 RocketMQProperties 配置类 +rocketmq: + name-server: 192.168.81.193:9876 + producer: + group: lms-producer-dev-group # 事务消息需要配置一样 + send-message-timeout: 3000 + consumer: + lms-task-operate: + group: lms-task-status-change-dev-group + topic: lms-task-status-change-dev-topic \ No newline at end of file diff --git a/nl-module-lms/nl-module-lms-server/src/main/resources/application.yaml b/nl-module-lms/nl-module-lms-server/src/main/resources/application.yaml index aa71aa8b..d811bc01 100644 --- a/nl-module-lms/nl-module-lms-server/src/main/resources/application.yaml +++ b/nl-module-lms/nl-module-lms-server/src/main/resources/application.yaml @@ -13,6 +13,7 @@ spring: import: - optional:classpath:application-${spring.profiles.active}.yaml # 加载【本地】配置 - optional:nacos:${spring.application.name}-${spring.profiles.active}.yaml # 加载【Nacos】的配置 + - optional:classpath:application-mq-${spring.profiles.active}.yaml # 加载 MQ 配置 # Servlet 配置 servlet: 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 2c5dae74..77f232d6 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 @@ -6,19 +6,21 @@ import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.springframework.stereotype.Component; import jakarta.annotation.Resource; +import org.springframework.beans.factory.annotation.Value; import java.util.UUID; /** * 任务事件 MQ 生产者 *

- * Topic: TASK_EVENT_TOPIC + * Topic 从 YAML 配置 rocketmq.consumer.task.topic 读取 * Tag: ownerService (LMS / WMS) */ @Slf4j @Component public class TaskEventProducer { - private static final String TOPIC = "TASK_EVENT_TOPIC"; + @Value("${rocketmq.consumer.task.topic}") + private String topic; @Resource private RocketMQTemplate rocketMQTemplate; @@ -34,7 +36,7 @@ public class TaskEventProducer { message.setEventId(UUID.randomUUID().toString()); } - String destination = TOPIC + "_" + message.getOwnerService(); + String destination = message.getOwnerService() + "_" + topic; try { rocketMQTemplate.syncSend(destination, message); log.info("MQ 发送成功, destination={}, taskId={}, eventType={}", 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 faf8d89f..49d41b9c 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 @@ -78,9 +78,12 @@ spring: # rocketmq 配置项,对应 RocketMQProperties 配置类 rocketmq: - name-server: 127.0.0.1:9876 # RocketMQ Namesrv + name-server: 192.168.81.193:9876 # RocketMQ Namesrv producer: - group: ${spring.application.name}_TASK_DEV_PRODUCER # 生产者分组 + group: task-producer-dev-group # 生产者分组 + consumer: + task: + topic: task-status-change-dev-topic spring: # RabbitMQ 配置项,对应 RabbitProperties 配置类 diff --git a/nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/TaskStatusChangeConsumer.java b/nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java similarity index 95% rename from nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/TaskStatusChangeConsumer.java rename to nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java index 57cbb940..fd51edee 100644 --- a/nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/TaskStatusChangeConsumer.java +++ b/nl-module-wms/nl-module-wms-server/src/main/java/cn/code/nl/module/wms/mq/consumer/WmsTaskStatusChangeConsumer.java @@ -22,7 +22,7 @@ import org.springframework.stereotype.Component; topic = "${rocketmq.consumer.wms-task-operate.topic}", consumerGroup = "${rocketmq.consumer.wms-task-operate.group}" ) -public class TaskStatusChangeConsumer implements RocketMQListener { +public class WmsTaskStatusChangeConsumer implements RocketMQListener { @Resource private TaskFactory taskFactory; diff --git a/nl-module-wms/nl-module-wms-server/src/main/resources/application-mq-dev.yaml b/nl-module-wms/nl-module-wms-server/src/main/resources/application-mq-dev.yaml index b0a11c37..325cbd30 100644 --- a/nl-module-wms/nl-module-wms-server/src/main/resources/application-mq-dev.yaml +++ b/nl-module-wms/nl-module-wms-server/src/main/resources/application-mq-dev.yaml @@ -2,6 +2,9 @@ # rocketmq 配置项,对应 RocketMQProperties 配置类 rocketmq: name-server: 192.168.81.193:9876 + producer: + group: wms-producer-dev-group # 事务消息需要配置一样 + send-message-timeout: 3000 consumer: wms-task-operate: group: wms-task-status-change-dev-group diff --git a/nl-server/src/main/resources/application-mq-dev.yaml b/nl-server/src/main/resources/application-mq-dev.yaml index b0a11c37..72ed6218 100644 --- a/nl-server/src/main/resources/application-mq-dev.yaml +++ b/nl-server/src/main/resources/application-mq-dev.yaml @@ -2,6 +2,9 @@ # rocketmq 配置项,对应 RocketMQProperties 配置类 rocketmq: name-server: 192.168.81.193:9876 + producer: + group: nl-producer-dev-group + send-message-timeout: 3000 consumer: wms-task-operate: group: wms-task-status-change-dev-group