askFileSys(@RequestBody AskFileSysMessage message){
String methodDescribe = getMethodDescribe("askFileSys");
- SendResult sendResult = askFileSysMessaggeTemplate.sendMember(message);
+ 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 0139672..2b9df9f 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,10 +7,13 @@ 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【需求编号】
@@ -19,16 +22,14 @@ import org.springframework.stereotype.Component;
* @version V1.0.0
*/
@Component
-public class AskFileSysMessaggeTemplate extends RocketMQEnhanceTemplate {
- public AskFileSysMessaggeTemplate(RocketMQTemplate template) {
- super(template);
- }
-
- public SendResult sendMember(AskFileSysMessage askRealDataMessage) {
+public class AskFileSysMessaggeTemplate {
+ @Resource
+ private MqTemplate mqTemplate;
+ public void sendMember(AskFileSysMessage askRealDataMessage) {
BaseMessage baseMessage = new BaseMessage();
baseMessage.setMessageBody(JSONObject.toJSONString(askRealDataMessage));
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
baseMessage.setKey(askRealDataMessage.getGuid());
- return send(BusinessTopic.FILE_TOPIC,askRealDataMessage.getNodeId() , baseMessage);
+ mqTemplate.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 18745af..1773c9d 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,10 +7,13 @@ 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【需求编号】
@@ -19,15 +22,13 @@ import org.springframework.stereotype.Component;
* @version V1.0.0
*/
@Component
-public class AskRealDataMessaggeTemplate extends RocketMQEnhanceTemplate {
- public AskRealDataMessaggeTemplate(RocketMQTemplate template) {
- super(template);
- }
-
- public SendResult sendMember(BaseMessage askRealDataMessage,String nodeId) {
+public class AskRealDataMessaggeTemplate {
+ @Resource
+ private MqTemplate mqTemplate;
+ public void sendMember(BaseMessage askRealDataMessage,String nodeId) {
askRealDataMessage.setSource(BusinessResource.WEB_RESOURCE);
AskRealDataMessage dto = JSON.parseObject(askRealDataMessage.getMessageBody(), AskRealDataMessage.class);
askRealDataMessage.setKey(dto.getLine());
- return send(BusinessTopic.ASK_REAL_DATA_TOPIC,nodeId , askRealDataMessage);
+ mqTemplate.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 d5b9338..bc5a413 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,10 +7,13 @@ 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;
+
/**
* 类的介绍:
*
@@ -19,14 +22,12 @@ import org.springframework.stereotype.Component;
* @createTime 2023/8/11 15:28
*/
@Component
-public class DeviceRebootMessageTemplate extends RocketMQEnhanceTemplate {
+public class DeviceRebootMessageTemplate {
- public DeviceRebootMessageTemplate(RocketMQTemplate template) {
- super(template);
- }
-
- public SendResult sendMember(BaseMessage baseMessage,String nodeId) {
+ @Resource
+ private MqTemplate mqTemplate;
+ public void sendMember(BaseMessage baseMessage,String nodeId) {
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
- return send(BusinessTopic.CONTROL_TOPIC,nodeId, baseMessage);
+ mqTemplate.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 8645108..3880a9b 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,10 +7,13 @@ 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;
+
/**
* 类的介绍:
*
@@ -19,16 +22,14 @@ import org.springframework.stereotype.Component;
* @createTime 2023/8/11 15:28
*/
@Component
-public class ProcessRebootMessageTemplate extends RocketMQEnhanceTemplate {
+public class ProcessRebootMessageTemplate {
- public ProcessRebootMessageTemplate(RocketMQTemplate template) {
- super(template);
- }
-
- public SendResult sendMember(BaseMessage baseMessage,String nodeId) {
+ @Resource
+ private MqTemplate mqTemplate;
+ public void sendMember(BaseMessage baseMessage,String nodeId) {
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
ProcessRebootMessage dto = JSON.parseObject(baseMessage.getMessageBody(), ProcessRebootMessage.class);
baseMessage.setKey(dto.getIndex()+"");
- return send(BusinessTopic.PROCESS_TOPIC,nodeId, baseMessage);
+ mqTemplate.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 f05d3d4..b3f7b9b 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,10 +5,13 @@ 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【需求编号】
@@ -17,13 +20,11 @@ import org.springframework.stereotype.Component;
* @version V1.0.0
*/
@Component
-public class RecallMessaggeTemplate extends RocketMQEnhanceTemplate {
- public RecallMessaggeTemplate(RocketMQTemplate template) {
- super(template);
- }
-
- public SendResult sendMember(BaseMessage recallMessage,String nodeId) {
+public class RecallMessaggeTemplate {
+ @Resource
+ private MqTemplate mqTemplate;
+ public void sendMember(BaseMessage recallMessage,String nodeId) {
recallMessage.setSource(BusinessResource.WEB_RESOURCE);
- return send(BusinessTopic.RECALL_TOPIC,nodeId , recallMessage);
+ mqTemplate.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 f724f80..239a2df 100644
--- a/message/message-boot/src/main/resources/bootstrap.yml
+++ b/message/message-boot/src/main/resources/bootstrap.yml
@@ -1,3 +1,5 @@
spring:
profiles:
- active: @spring.profiles.active@
\ No newline at end of file
+ active: @spring.profiles.active@
+mq:
+ type: redis-stream
\ 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 de2d75f..bab4add 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