askFileSys(@RequestBody AskFileSysMessage message){
String methodDescribe = getMethodDescribe("askFileSys");
- askFileSysMessaggeTemplate.sendMember(message);
+ SendResult sendResult = askFileSysMessaggeTemplate.sendMember(message);
return HttpResultUtil.assembleCommonResponseResult(CommonResponseEnum.SUCCESS,message.getGuid() , methodDescribe);
}
}
diff --git a/message/message-boot/src/main/java/com/njcn/message/produce/template/AskFileSysMessaggeTemplate.java b/message/message-boot/src/main/java/com/njcn/message/produce/template/AskFileSysMessaggeTemplate.java
index 2b9df9f..0139672 100644
--- a/message/message-boot/src/main/java/com/njcn/message/produce/template/AskFileSysMessaggeTemplate.java
+++ b/message/message-boot/src/main/java/com/njcn/message/produce/template/AskFileSysMessaggeTemplate.java
@@ -7,13 +7,10 @@ import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.AskFileSysMessage;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
-import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
-import javax.annotation.Resource;
-
/**
* Description:
* Date: 2024/12/13 15:15【需求编号】
@@ -22,14 +19,16 @@ import javax.annotation.Resource;
* @version V1.0.0
*/
@Component
-public class AskFileSysMessaggeTemplate {
- @Resource
- private MqTemplate mqTemplate;
- public void sendMember(AskFileSysMessage askRealDataMessage) {
+public class AskFileSysMessaggeTemplate extends RocketMQEnhanceTemplate {
+ public AskFileSysMessaggeTemplate(RocketMQTemplate template) {
+ super(template);
+ }
+
+ public SendResult sendMember(AskFileSysMessage askRealDataMessage) {
BaseMessage baseMessage = new BaseMessage();
baseMessage.setMessageBody(JSONObject.toJSONString(askRealDataMessage));
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
baseMessage.setKey(askRealDataMessage.getGuid());
- mqTemplate.send(BusinessTopic.FILE_TOPIC,askRealDataMessage.getNodeId() , baseMessage);
+ return send(BusinessTopic.FILE_TOPIC,askRealDataMessage.getNodeId() , baseMessage);
}
}
diff --git a/message/message-boot/src/main/java/com/njcn/message/produce/template/AskRealDataMessaggeTemplate.java b/message/message-boot/src/main/java/com/njcn/message/produce/template/AskRealDataMessaggeTemplate.java
index 1773c9d..18745af 100644
--- a/message/message-boot/src/main/java/com/njcn/message/produce/template/AskRealDataMessaggeTemplate.java
+++ b/message/message-boot/src/main/java/com/njcn/message/produce/template/AskRealDataMessaggeTemplate.java
@@ -7,13 +7,10 @@ import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.AskRealDataMessage;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
-import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
-import javax.annotation.Resource;
-
/**
* Description:
* Date: 2024/12/13 15:15【需求编号】
@@ -22,13 +19,15 @@ import javax.annotation.Resource;
* @version V1.0.0
*/
@Component
-public class AskRealDataMessaggeTemplate {
- @Resource
- private MqTemplate mqTemplate;
- public void sendMember(BaseMessage askRealDataMessage,String nodeId) {
+public class AskRealDataMessaggeTemplate extends RocketMQEnhanceTemplate {
+ public AskRealDataMessaggeTemplate(RocketMQTemplate template) {
+ super(template);
+ }
+
+ public SendResult sendMember(BaseMessage askRealDataMessage,String nodeId) {
askRealDataMessage.setSource(BusinessResource.WEB_RESOURCE);
AskRealDataMessage dto = JSON.parseObject(askRealDataMessage.getMessageBody(), AskRealDataMessage.class);
askRealDataMessage.setKey(dto.getLine());
- mqTemplate.send(BusinessTopic.ASK_REAL_DATA_TOPIC,nodeId , askRealDataMessage);
+ return send(BusinessTopic.ASK_REAL_DATA_TOPIC,nodeId , askRealDataMessage);
}
}
diff --git a/message/message-boot/src/main/java/com/njcn/message/produce/template/DeviceRebootMessageTemplate.java b/message/message-boot/src/main/java/com/njcn/message/produce/template/DeviceRebootMessageTemplate.java
index bc5a413..d5b9338 100644
--- a/message/message-boot/src/main/java/com/njcn/message/produce/template/DeviceRebootMessageTemplate.java
+++ b/message/message-boot/src/main/java/com/njcn/message/produce/template/DeviceRebootMessageTemplate.java
@@ -7,13 +7,10 @@ import com.njcn.message.message.AskRealDataMessage;
import com.njcn.message.message.DeviceRebootMessage;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
-import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
-import javax.annotation.Resource;
-
/**
* 类的介绍:
*
@@ -22,12 +19,14 @@ import javax.annotation.Resource;
* @createTime 2023/8/11 15:28
*/
@Component
-public class DeviceRebootMessageTemplate {
+public class DeviceRebootMessageTemplate extends RocketMQEnhanceTemplate {
- @Resource
- private MqTemplate mqTemplate;
- public void sendMember(BaseMessage baseMessage,String nodeId) {
+ public DeviceRebootMessageTemplate(RocketMQTemplate template) {
+ super(template);
+ }
+
+ public SendResult sendMember(BaseMessage baseMessage,String nodeId) {
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
- mqTemplate.send(BusinessTopic.CONTROL_TOPIC,nodeId, baseMessage);
+ return send(BusinessTopic.CONTROL_TOPIC,nodeId, baseMessage);
}
}
diff --git a/message/message-boot/src/main/java/com/njcn/message/produce/template/ProcessRebootMessageTemplate.java b/message/message-boot/src/main/java/com/njcn/message/produce/template/ProcessRebootMessageTemplate.java
index 3880a9b..8645108 100644
--- a/message/message-boot/src/main/java/com/njcn/message/produce/template/ProcessRebootMessageTemplate.java
+++ b/message/message-boot/src/main/java/com/njcn/message/produce/template/ProcessRebootMessageTemplate.java
@@ -7,13 +7,10 @@ import com.njcn.message.message.AskRealDataMessage;
import com.njcn.message.message.ProcessRebootMessage;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
-import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
-import javax.annotation.Resource;
-
/**
* 类的介绍:
*
@@ -22,14 +19,16 @@ import javax.annotation.Resource;
* @createTime 2023/8/11 15:28
*/
@Component
-public class ProcessRebootMessageTemplate {
+public class ProcessRebootMessageTemplate extends RocketMQEnhanceTemplate {
- @Resource
- private MqTemplate mqTemplate;
- public void sendMember(BaseMessage baseMessage,String nodeId) {
+ public ProcessRebootMessageTemplate(RocketMQTemplate template) {
+ super(template);
+ }
+
+ public SendResult sendMember(BaseMessage baseMessage,String nodeId) {
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
ProcessRebootMessage dto = JSON.parseObject(baseMessage.getMessageBody(), ProcessRebootMessage.class);
baseMessage.setKey(dto.getIndex()+"");
- mqTemplate.send(BusinessTopic.PROCESS_TOPIC,nodeId, baseMessage);
+ return send(BusinessTopic.PROCESS_TOPIC,nodeId, baseMessage);
}
}
diff --git a/message/message-boot/src/main/java/com/njcn/message/produce/template/RecallMessaggeTemplate.java b/message/message-boot/src/main/java/com/njcn/message/produce/template/RecallMessaggeTemplate.java
index b3f7b9b..f05d3d4 100644
--- a/message/message-boot/src/main/java/com/njcn/message/produce/template/RecallMessaggeTemplate.java
+++ b/message/message-boot/src/main/java/com/njcn/message/produce/template/RecallMessaggeTemplate.java
@@ -5,13 +5,10 @@ import com.njcn.message.constant.BusinessResource;
import com.njcn.message.constant.BusinessTopic;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
-import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
-import javax.annotation.Resource;
-
/**
* Description:
* Date: 2024/12/13 15:15【需求编号】
@@ -20,11 +17,13 @@ import javax.annotation.Resource;
* @version V1.0.0
*/
@Component
-public class RecallMessaggeTemplate {
- @Resource
- private MqTemplate mqTemplate;
- public void sendMember(BaseMessage recallMessage,String nodeId) {
+public class RecallMessaggeTemplate extends RocketMQEnhanceTemplate {
+ public RecallMessaggeTemplate(RocketMQTemplate template) {
+ super(template);
+ }
+
+ public SendResult sendMember(BaseMessage recallMessage,String nodeId) {
recallMessage.setSource(BusinessResource.WEB_RESOURCE);
- mqTemplate.send(BusinessTopic.RECALL_TOPIC,nodeId , recallMessage);
+ return send(BusinessTopic.RECALL_TOPIC,nodeId , recallMessage);
}
}
diff --git a/message/message-boot/src/main/resources/bootstrap.yml b/message/message-boot/src/main/resources/bootstrap.yml
index 239a2df..f724f80 100644
--- a/message/message-boot/src/main/resources/bootstrap.yml
+++ b/message/message-boot/src/main/resources/bootstrap.yml
@@ -1,5 +1,3 @@
spring:
profiles:
- active: @spring.profiles.active@
-mq:
- type: redis-stream
\ No newline at end of file
+ active: @spring.profiles.active@
\ No newline at end of file
diff --git a/message/message-boot/src/test/java/com/njcn/message/mq/MqSliceWiringTest.java b/message/message-boot/src/test/java/com/njcn/message/mq/MqSliceWiringTest.java
index bab4add..de2d75f 100644
--- a/message/message-boot/src/test/java/com/njcn/message/mq/MqSliceWiringTest.java
+++ b/message/message-boot/src/test/java/com/njcn/message/mq/MqSliceWiringTest.java
@@ -1,210 +1,210 @@
-//package com.njcn.message.mq;
-//
-//import com.alibaba.fastjson.JSON;
-//import com.njcn.message.message.RecallMessage;
-//import com.njcn.message.messagedto.MessageDataDTO;
-//import com.njcn.middle.stream.autoconfig.RedisStreamAutoConfiguration;
-//import com.njcn.mq.autoconfig.MqCoreAutoConfiguration;
-//import com.njcn.mq.container.MqListenerRegistry;
-//import com.njcn.mq.core.MqTemplate;
-//import com.njcn.mq.driver.redis.RedisStreamMqDriver;
-//import com.njcn.mq.driver.redis.RedisStreamMqDriverAutoConfiguration;
-//import com.njcn.mq.driver.rocketmq.RocketMqDriver;
-//import com.njcn.mq.driver.rocketmq.RocketMqDriverAutoConfiguration;
-//import com.njcn.mq.spi.MqDriver;
-//import com.njcn.stat.api.MessAnalysisFeignClient;
-//import org.apache.rocketmq.spring.autoconfigure.RocketMQProperties;
-//import org.apache.rocketmq.spring.core.RocketMQTemplate;
-//import org.junit.jupiter.api.Assumptions;
-//import org.junit.jupiter.api.Test;
-//import org.springframework.boot.autoconfigure.AutoConfigurations;
-//import org.springframework.boot.autoconfigure.data.redis.RedisAutoConfiguration;
-//import org.springframework.boot.test.context.runner.ApplicationContextRunner;
-//import org.springframework.data.redis.connection.stream.MapRecord;
-//import org.springframework.data.redis.connection.stream.ReadOffset;
-//import org.springframework.data.redis.connection.stream.StreamOffset;
-//import org.springframework.data.redis.connection.stream.StreamReadOptions;
-//import org.springframework.data.redis.connection.stream.StreamRecords;
-//import org.springframework.data.redis.core.StringRedisTemplate;
-//
-//import java.net.InetSocketAddress;
-//import java.net.Socket;
-//import java.util.HashMap;
-//import java.util.List;
-//import java.util.Map;
-//
-//import static org.assertj.core.api.Assertions.assertThat;
-//import static org.mockito.ArgumentMatchers.anyList;
-//import static org.mockito.Mockito.mock;
-//import static org.mockito.Mockito.timeout;
-//import static org.mockito.Mockito.verify;
-//
-///**
-// * MQ starter 切片接入验证测试(ApplicationContextRunner,不启动全应用)。
-// * 真实 Redis 要求:localhost:6379(docker redis-stream-it);仅上/下行收发用例需要,
-// * 纯装配用例不需要(Lettuce 懒连接,不触发实际连接)。
-// *
-// * 装配一律用 {@code withConfiguration(AutoConfigurations.of(...))},让 Spring 按
-// * 真实 autoconfig 顺序(AutoConfigurationSorter)处理,与生产环境一致。以此
-// * 验证 {@link MqCoreAutoConfiguration} 的 {@code @ConditionalOnBean(MqDriver)} 能在
-// * 真实顺序下拿到 driver bean —— 这依赖 driver autoconfig 上的
-// * {@code @AutoConfigureBefore(MqCoreAutoConfiguration.class)}。若缺该保障,按类名字母序
-// * {@code com.njcn.mq.autoconfig.*} 先于 {@code com.njcn.mq.driver.*},core 会在 driver
-// * bean 注册前求值条件失败,导致 MqTemplate / MqListenerRegistry 不装配。
-// */
-//public class MqSliceWiringTest {
-//
-// private static boolean redisAvailable() {
-// try {
-// Socket s = new Socket();
-// s.connect(new InetSocketAddress("localhost", 6379), 500);
-// s.close();
-// return true;
-// } catch (Exception e) {
-// return false;
-// }
-// }
-//
-// /** redis-stream 模式:真实 autoconfig 顺序装 5 个 autoconfig + 业务 bean(上/下行用)。 */
-// private ApplicationContextRunner redisStreamRunner() {
-// return new ApplicationContextRunner()
-// .withConfiguration(AutoConfigurations.of(
-// RedisAutoConfiguration.class,
-// RedisStreamAutoConfiguration.class,
-// RedisStreamMqDriverAutoConfiguration.class,
-// RocketMqDriverAutoConfiguration.class,
-// MqCoreAutoConfiguration.class))
-// .withUserConfiguration(FrontDataMqListener.class, RecallMqSender.class)
-// .withBean(MessAnalysisFeignClient.class,
-// () -> mock(MessAnalysisFeignClient.class))
-// .withPropertyValues(
-// "mq.type=redis-stream",
-// "spring.redis.host=localhost",
-// "spring.redis.port=6379");
-// }
-//
-// // ------------------------------------------------------------------
-// // 1. 装配:mq.type 二选一 + core bean 在真实 autoconfig 顺序下装配
-// // ------------------------------------------------------------------
-//
-// @Test
-// void testRedisStreamDriverSelected() {
-// // 纯装配检查,不依赖真实 Redis(Lettuce 懒连接)。
-// new ApplicationContextRunner()
-// .withConfiguration(AutoConfigurations.of(
-// RedisAutoConfiguration.class,
-// RedisStreamAutoConfiguration.class,
-// RedisStreamMqDriverAutoConfiguration.class,
-// RocketMqDriverAutoConfiguration.class,
-// MqCoreAutoConfiguration.class))
-// .withPropertyValues(
-// "mq.type=redis-stream",
-// "spring.redis.host=localhost",
-// "spring.redis.port=6379")
-// .run(ctx -> {
-// assertThat(ctx).hasSingleBean(MqDriver.class);
-// assertThat(ctx.getBean(MqDriver.class)).isInstanceOf(RedisStreamMqDriver.class);
-// // 真实顺序下 core 的 @ConditionalOnBean(MqDriver) 必须命中
-// assertThat(ctx).hasSingleBean(MqTemplate.class);
-// assertThat(ctx).hasSingleBean(MqListenerRegistry.class);
-// });
-// }
-//
-// @Test
-// void testRocketMqDriverSelected() {
-// // rocket 测试不依赖真实 Redis;不装 RedisAutoConfiguration,redis driver 因 mq.type 不匹配而 inactive。
-// new ApplicationContextRunner()
-// .withConfiguration(AutoConfigurations.of(
-// RedisStreamMqDriverAutoConfiguration.class,
-// RocketMqDriverAutoConfiguration.class,
-// MqCoreAutoConfiguration.class))
-// .withBean(RocketMQTemplate.class, () -> mock(RocketMQTemplate.class))
-// .withBean(RocketMQProperties.class, RocketMQProperties::new)
-// .withPropertyValues("mq.type=rocketmq")
-// .run(ctx -> {
-// assertThat(ctx).hasSingleBean(MqDriver.class);
-// assertThat(ctx.getBean(MqDriver.class)).isInstanceOf(RocketMqDriver.class);
-// assertThat(ctx).hasSingleBean(MqTemplate.class);
-// assertThat(ctx).hasSingleBean(MqListenerRegistry.class);
-// });
-// }
-//
-// // ------------------------------------------------------------------
-// // 2. 上行(需要真实 Redis)
-// // 向 LN_Topic stream XADD 一条消息,验证 FrontDataMqListener 触发了
-// // messAnalysisFeignClient.analysis(...)
-// // ------------------------------------------------------------------
-//
-// @Test
-// void testUpstreamFrontData() {
-// Assumptions.assumeTrue(redisAvailable(), "Redis not available at localhost:6379, skipping");
-// redisStreamRunner().run(ctx -> {
-// StringRedisTemplate redis = ctx.getBean(StringRedisTemplate.class);
-// MessAnalysisFeignClient mockClient = ctx.getBean(MessAnalysisFeignClient.class);
-//
-// // 构造 RedisStreamMqDriver 消费时期望的字段格式
-// // enc=plain → body 即为原始 JSON(参见 MessageCodec.encodeBody(json, false))
-// // stream key = StreamKeyBuilder.buildStream("LN_Topic") = "LN_Topic"
-// // (默认 envIsolation=false,无 env 前缀)
-// MessageDataDTO dto = new MessageDataDTO();
-// dto.setDataType(1);
-// String json = JSON.toJSONString(dto);
-//
-// Map fields = new HashMap();
-// fields.put("key", "test-key-upstream-1");
-// fields.put("tag", "");
-// fields.put("enc", "plain");
-// fields.put("body", json);
-// fields.put("ts", String.valueOf(System.currentTimeMillis()));
-//
-// // XADD 到 LN_Topic(MqListenerRegistry.start() 已在 context 启动时订阅该 stream)
-// redis.opsForStream().add(
-// StreamRecords.newRecord().in("LN_Topic").ofMap(fields)
-// );
-//
-// // 等待消费线程处理(blockMs 默认 2000ms,timeout=5000ms 足够)
-// verify(mockClient, timeout(5000)).analysis(anyList());
-// });
-// }
-//
-// // ------------------------------------------------------------------
-// // 3. 下行(需要真实 Redis)
-// // 通过 RecallMqSender 发送 RecallMessage,验证 stream 内容正确
-// // ------------------------------------------------------------------
-//
-// @Test
-// void testDownstreamRecall() {
-// Assumptions.assumeTrue(redisAvailable(), "Redis not available at localhost:6379, skipping");
-// redisStreamRunner().run(ctx -> {
-// StringRedisTemplate redis = ctx.getBean(StringRedisTemplate.class);
-// RecallMqSender sender = ctx.getBean(RecallMqSender.class);
-//
-// // 清理,确保只有本次发送的消息
-// redis.delete("n1_recall_Topic");
-//
-// RecallMessage msg = new RecallMessage();
-// msg.setGuid("g1");
-// // MqTemplate.send("n1_recall_Topic", msg) → RedisStreamMqDriver XADD
-// // stream key = StreamKeyBuilder.buildStream("n1_recall_Topic") = "n1_recall_Topic"
-// sender.send(msg, "n1");
-//
-// Long size = redis.opsForStream().size("n1_recall_Topic");
-// assertThat(size).isGreaterThanOrEqualTo(1L);
-//
-// @SuppressWarnings("unchecked")
-// List> records =
-// redis.opsForStream().read(
-// StreamReadOptions.empty().count(1L),
-// StreamOffset.create("n1_recall_Topic", ReadOffset.from("0-0"))
-// );
-// assertThat(records).isNotEmpty();
-//
-// Map