1 Commits

Author SHA1 Message Date
hzj
e982fa960e 冀北版本适配 2026-05-09 08:43:57 +08:00
10 changed files with 22 additions and 14 deletions

View File

@@ -33,7 +33,7 @@ import java.util.Objects;
@RocketMQMessageListener( @RocketMQMessageListener(
topic = "Device_Run_Flag_Topic", topic = "Device_Run_Flag_Topic",
consumerGroup = "Device_Run_Flag_Consumer", consumerGroup = "Device_Run_Flag_Consumer",
selectorExpression = "Test_Tag||Test_Keys", selectorExpression = "*",
consumeThreadNumber = 10, consumeThreadNumber = 10,
enableMsgTrace = true enableMsgTrace = true
) )

View File

@@ -42,7 +42,7 @@ import java.util.concurrent.TimeUnit;
@RocketMQMessageListener( @RocketMQMessageListener(
topic = "LN_Topic", topic = "LN_Topic",
consumerGroup = "ln_consumer", consumerGroup = "ln_consumer",
selectorExpression = "Test_Tag||Test_Keys", selectorExpression = "*",
consumeThreadNumber = 10, consumeThreadNumber = 10,
enableMsgTrace = true enableMsgTrace = true
) )

View File

@@ -33,7 +33,7 @@ import java.util.Objects;
@RocketMQMessageListener( @RocketMQMessageListener(
topic = "Heart_Beat_Topic", topic = "Heart_Beat_Topic",
consumerGroup = "Heartb_Beat_Topic_Consumer", consumerGroup = "Heartb_Beat_Topic_Consumer",
selectorExpression = "Test_Tag||Test_Keys", selectorExpression = "*",
consumeThreadNumber = 10, consumeThreadNumber = 10,
enableMsgTrace = true enableMsgTrace = true
) )

View File

@@ -31,7 +31,7 @@ import java.util.Objects;
@RocketMQMessageListener( @RocketMQMessageListener(
topic = "Real_Time_Data_Topic", topic = "Real_Time_Data_Topic",
consumerGroup = "real_time_consumer", consumerGroup = "real_time_consumer",
selectorExpression = "Test_Tag||Test_Keys", selectorExpression = "*",
consumeThreadNumber = 10, consumeThreadNumber = 10,
enableMsgTrace = true enableMsgTrace = true
) )

View File

@@ -33,7 +33,7 @@ import java.util.Objects;
@RocketMQMessageListener( @RocketMQMessageListener(
topic = "log_Topic", topic = "log_Topic",
consumerGroup = "Log_Topic_Consumer", consumerGroup = "Log_Topic_Consumer",
selectorExpression = "Test_Tag||Test_Keys", selectorExpression = "*",
consumeThreadNumber = 10, consumeThreadNumber = 10,
enableMsgTrace = true enableMsgTrace = true
) )

View File

@@ -32,7 +32,7 @@ import java.util.Objects;
@RocketMQMessageListener( @RocketMQMessageListener(
topic = "Topic_Reply_Topic", topic = "Topic_Reply_Topic",
consumerGroup = "Topic_Reply_Topic_Consumer", consumerGroup = "Topic_Reply_Topic_Consumer",
selectorExpression = "Test_Tag||Test_Keys", selectorExpression = "*",
consumeThreadNumber = 10, consumeThreadNumber = 10,
enableMsgTrace = true enableMsgTrace = true
) )

View File

@@ -1,8 +1,10 @@
package com.njcn.message.produce.template; package com.njcn.message.produce.template;
import com.alibaba.fastjson.JSON;
import com.njcn.message.constant.BusinessResource; import com.njcn.message.constant.BusinessResource;
import com.njcn.message.constant.BusinessTopic; import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.AskRealDataMessage;
import com.njcn.middle.rocket.domain.BaseMessage; import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate; import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendResult;
@@ -24,7 +26,8 @@ public class AskRealDataMessaggeTemplate extends RocketMQEnhanceTemplate {
public SendResult sendMember(BaseMessage askRealDataMessage,String nodeId) { public SendResult sendMember(BaseMessage askRealDataMessage,String nodeId) {
askRealDataMessage.setSource(BusinessResource.WEB_RESOURCE); askRealDataMessage.setSource(BusinessResource.WEB_RESOURCE);
askRealDataMessage.setKey("Test_Keys"); AskRealDataMessage dto = JSON.parseObject(askRealDataMessage.getMessageBody(), AskRealDataMessage.class);
return send(nodeId+"_"+BusinessTopic.ASK_REAL_DATA_TOPIC,"Test_Tag" , askRealDataMessage); askRealDataMessage.setKey(dto.getLine());
return send(BusinessTopic.ASK_REAL_DATA_TOPIC,nodeId , askRealDataMessage);
} }
} }

View File

@@ -1,7 +1,10 @@
package com.njcn.message.produce.template; package com.njcn.message.produce.template;
import com.alibaba.fastjson.JSON;
import com.njcn.message.constant.BusinessResource; import com.njcn.message.constant.BusinessResource;
import com.njcn.message.constant.BusinessTopic; import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.AskRealDataMessage;
import com.njcn.message.message.DeviceRebootMessage;
import com.njcn.middle.rocket.domain.BaseMessage; import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate; import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendResult;
@@ -24,7 +27,6 @@ public class DeviceRebootMessageTemplate extends RocketMQEnhanceTemplate {
public SendResult sendMember(BaseMessage baseMessage,String nodeId) { public SendResult sendMember(BaseMessage baseMessage,String nodeId) {
baseMessage.setSource(BusinessResource.WEB_RESOURCE); baseMessage.setSource(BusinessResource.WEB_RESOURCE);
baseMessage.setKey("Test_Keys"); return send(BusinessTopic.CONTROL_TOPIC,nodeId, baseMessage);
return send(nodeId+"_"+BusinessTopic.CONTROL_TOPIC,"Test_Tag", baseMessage);
} }
} }

View File

@@ -1,7 +1,10 @@
package com.njcn.message.produce.template; package com.njcn.message.produce.template;
import com.alibaba.fastjson.JSON;
import com.njcn.message.constant.BusinessResource; import com.njcn.message.constant.BusinessResource;
import com.njcn.message.constant.BusinessTopic; import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.AskRealDataMessage;
import com.njcn.message.message.ProcessRebootMessage;
import com.njcn.middle.rocket.domain.BaseMessage; import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate; import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendResult;
@@ -24,7 +27,8 @@ public class ProcessRebootMessageTemplate extends RocketMQEnhanceTemplate {
public SendResult sendMember(BaseMessage baseMessage,String nodeId) { public SendResult sendMember(BaseMessage baseMessage,String nodeId) {
baseMessage.setSource(BusinessResource.WEB_RESOURCE); baseMessage.setSource(BusinessResource.WEB_RESOURCE);
baseMessage.setKey("Test_Keys"); ProcessRebootMessage dto = JSON.parseObject(baseMessage.getMessageBody(), ProcessRebootMessage.class);
return send(nodeId+"_"+BusinessTopic.PROCESS_TOPIC,"Test_Tag", baseMessage); baseMessage.setKey(dto.getIndex()+"");
return send(BusinessTopic.PROCESS_TOPIC,nodeId, baseMessage);
} }
} }

View File

@@ -24,7 +24,6 @@ public class RecallMessaggeTemplate extends RocketMQEnhanceTemplate {
public SendResult sendMember(BaseMessage recallMessage,String nodeId) { public SendResult sendMember(BaseMessage recallMessage,String nodeId) {
recallMessage.setSource(BusinessResource.WEB_RESOURCE); recallMessage.setSource(BusinessResource.WEB_RESOURCE);
recallMessage.setKey("Test_Keys"); return send(BusinessTopic.RECALL_TOPIC,nodeId , recallMessage);
return send(nodeId+"_"+BusinessTopic.RECALL_TOPIC,"Test_Tag" , recallMessage);
} }
} }