优化交换机数据入库

This commit is contained in:
gaoyutao
2025-09-28 19:33:39 +08:00
parent 9fa62124e6
commit b8548b1077
21 changed files with 331 additions and 115 deletions
@@ -115,7 +115,7 @@ public class RocketMqController {
private Map sendDelayMessage(@RequestParam("topic") String topic,@RequestParam("tag") String tag,@RequestParam("key") String key,@RequestParam("value") String value){
MessageProducer messageProducer = new MessageProducer();
//调用MessageProducer配置好的消息方法 topic需要你根据你们业务定制相应的
SendResult sendResult = messageProducer.sendDelayMessage("order-message","order_timeout_tag","title","content");
SendResult sendResult = messageProducer.sendDelayMessage("order-message","order_timeout_tag","title","content",4);
Map<String,Object> result = new HashMap<>();
result.put("data",sendResult);
return result;
@@ -342,14 +342,15 @@ public class DeviceMessageHandler {
private void handleSwitchPwrMessage(CollectDataVo switchDataVo, String clientId){
List<InitialSwitchPowerSupply> powerSupplyList = SwitchJsonDataParser.parseJsonData(switchDataVo.getValue(), InitialSwitchPowerSupply.class);
if (!powerSupplyList.isEmpty()){
InitialSwitchPowerSupply insertData = powerSupplyList.get(0);
// 时间戳转换
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
initialSwitchPowerSupplyService.insertInitialSwitchPowerSupply(insertData);
for (InitialSwitchPowerSupply insertData : powerSupplyList) {
// 时间戳转换
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
}
initialSwitchPowerSupplyService.insertBatchInitialSwitchPowerSupply(powerSupplyList);
}
}
/**
@@ -359,14 +360,15 @@ public class DeviceMessageHandler {
private void handleSwitchModuleMessage(CollectDataVo switchDataVo, String clientId){
List<InitialSwitchOpticalModule> moduleList = SwitchJsonDataParser.parseJsonData(switchDataVo.getValue(), InitialSwitchOpticalModule.class);
if (!moduleList.isEmpty()){
InitialSwitchOpticalModule insertData = moduleList.get(0);
// 时间戳转换
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
initialSwitchOpticalModuleService.insertInitialSwitchOpticalModule(insertData);
for (InitialSwitchOpticalModule insertData : moduleList) {
// 时间戳转换
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
}
initialSwitchOpticalModuleService.batchInitialSwitchOpticalModule(moduleList);
}
}
@@ -377,14 +379,15 @@ public class DeviceMessageHandler {
private void handleSwitchMpuMessage(CollectDataVo switchDataVo, String clientId){
List<InitialSwitchMpuInfo> mpuList = SwitchJsonDataParser.parseJsonData(switchDataVo.getValue(), InitialSwitchMpuInfo.class);
if (!mpuList.isEmpty()){
InitialSwitchMpuInfo insertData = mpuList.get(0);
// 时间戳转换
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
initialSwitchMpuInfoService.insertInitialSwitchMpuInfo(insertData);
for (InitialSwitchMpuInfo insertData : mpuList) {
// 时间戳转换
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
}
initialSwitchMpuInfoService.insertBatchInitialSwitchMpuInfo(mpuList);
}
}
@@ -395,14 +398,15 @@ public class DeviceMessageHandler {
private void handleSwitchFanMessage(CollectDataVo switchDataVo, String clientId){
List<InitialSwitchFanInfo> fanList = SwitchJsonDataParser.parseJsonData(switchDataVo.getValue(), InitialSwitchFanInfo.class);
if (!fanList.isEmpty()){
InitialSwitchFanInfo insertData = fanList.get(0);
// 时间戳转换
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
initialSwitchFanInfoService.insertInitialSwitchFanInfo(insertData);
for (InitialSwitchFanInfo insertData : fanList) {
// 时间戳转换
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
}
initialSwitchFanInfoService.insertBatchInitialSwitchFanInfo(fanList);
}
}
@@ -410,54 +414,47 @@ public class DeviceMessageHandler {
* 其他发现数据(默认处理)
* @param switchDataVo
*/
private void handleSwitchOtherMessage(CollectDataVo switchDataVo, String clientId){
if (switchDataVo != null){
private void handleSwitchOtherMessage(CollectDataVo switchDataVo, String clientId) {
if (switchDataVo != null) {
try {
InitialSwitchOtherCollectData insertData = new InitialSwitchOtherCollectData();
// 时间戳转换
// 设置基本信息
long timestamp = switchDataVo.getTimestamp();
long millis = timestamp * 1000;
Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒
Date createTime = new Date(timestamp * 1000 / 1000 * 1000);
insertData.setClientId(clientId);
insertData.setCreateTime(createTime);
String value = switchDataVo.getValue();
// 解析JSON字符串
JsonNode jsonNode = objectMapper.readTree(value);
// 获取第一个字段的key和value
if (jsonNode.isObject() && jsonNode.fields().hasNext()) {
Map.Entry<String, JsonNode> entry = jsonNode.fields().next();
String fieldName = entry.getKey();
String fieldValue = entry.getValue().asText();
insertData.setCollectType(fieldName);
if(!"null".equals(fieldValue)){
insertData.setCollectValue(fieldValue);
}
} else if (jsonNode.isArray() && jsonNode.size() > 0) {
// 处理数组格式 [{}]
// 处理数组中的字符串JSON
if (jsonNode.isArray() && jsonNode.size() > 0) {
JsonNode firstElement = jsonNode.get(0);
if (firstElement.isObject() && firstElement.fields().hasNext()) {
Map.Entry<String, JsonNode> entry = firstElement.fields().next();
String fieldName = entry.getKey();
String fieldValue = entry.getValue().asText();
insertData.setCollectType(fieldName);
if(!"null".equals(fieldValue)){
insertData.setCollectValue(fieldValue);
if (firstElement.isTextual()) {
// 二次解析JSON字符串
JsonNode innerJsonNode = objectMapper.readTree(firstElement.asText());
if (innerJsonNode.isObject()) {
Iterator<Map.Entry<String, JsonNode>> fields = innerJsonNode.fields();
if (fields.hasNext()) {
Map.Entry<String, JsonNode> entry = fields.next();
String fieldName = entry.getKey();
String fieldValue = entry.getValue().asText();
insertData.setCollectType(fieldName);
if (!"null".equals(fieldValue)) {
insertData.setCollectValue(fieldValue);
insertInitialSwitchOtherInfo.insertInitialSwitchOtherCollectData(insertData);
}
}
}
}
}
insertInitialSwitchOtherInfo.insertInitialSwitchOtherCollectData(insertData);
} catch (Exception e) {
log.error("解析JSON数据失败: {}, value: {}", e.getMessage(), switchDataVo.getValue());
// 可以选择保存原始数据或进行其他错误处理
log.error("解析JSON数据失败: {}, value: {}", e.getMessage(), switchDataVo.getValue(), e);
}
}
}
@@ -1,8 +1,9 @@
package com.ruoyi.rocketmq.mapper;
import java.util.List;
import com.ruoyi.rocketmq.domain.InitialSwitchFanInfo;
import java.util.List;
/**
* 风扇信息Mapper接口
*
@@ -58,4 +59,11 @@ public interface InitialSwitchFanInfoMapper
* @return 结果
*/
public int deleteInitialSwitchFanInfoByIds(Long[] ids);
/**
* 批量新增风扇信息
* @param list
* @return
*/
public int insertBatchInitialSwitchFanInfo(List<InitialSwitchFanInfo> list);
}
@@ -1,8 +1,9 @@
package com.ruoyi.rocketmq.mapper;
import java.util.List;
import com.ruoyi.rocketmq.domain.InitialSwitchMpuInfo;
import java.util.List;
/**
* MPU信息Mapper接口
*
@@ -58,4 +59,11 @@ public interface InitialSwitchMpuInfoMapper
* @return 结果
*/
public int deleteInitialSwitchMpuInfoByIds(Long[] ids);
/**
* 批量新增mpu信息
* @param list
* @return
*/
int insertBatchInitialSwitchMpuInfo(List<InitialSwitchMpuInfo> list);
}
@@ -1,8 +1,9 @@
package com.ruoyi.rocketmq.mapper;
import java.util.List;
import com.ruoyi.rocketmq.domain.InitialSwitchOpticalModule;
import java.util.List;
/**
* 光模块信息Mapper接口
*
@@ -58,4 +59,10 @@ public interface InitialSwitchOpticalModuleMapper
* @return 结果
*/
public int deleteInitialSwitchOpticalModuleByIds(Long[] ids);
/**
* 批量新增光模块信息
* @param list
*/
int batchInitialSwitchOpticalModule(List<InitialSwitchOpticalModule> list);
}
@@ -1,8 +1,9 @@
package com.ruoyi.rocketmq.mapper;
import java.util.List;
import com.ruoyi.rocketmq.domain.InitialSwitchPowerSupply;
import java.util.List;
/**
* 电源信息Mapper接口
*
@@ -58,4 +59,10 @@ public interface InitialSwitchPowerSupplyMapper
* @return 结果
*/
public int deleteInitialSwitchPowerSupplyByIds(Long[] ids);
/**
* 批量新增电源信息
* @param list
* @return
*/
public int insertBatchInitialSwitchPowerSupply(List<InitialSwitchPowerSupply> list);
}
@@ -158,15 +158,16 @@ public class MessageProducer {
* @param value 消息的内容
* 延迟发送:通过设置延迟级别来实现延迟发送消息。
*/
public SendResult sendDelayMessage(String topic, String tag, String key, String value)
public SendResult sendDelayMessage(String topic, String tag, String key, String value, Integer level)
{
SendResult result = null;
try
{
Message msg = new Message(topic,tag,key, value.getBytes(RemotingHelper.DEFAULT_CHARSET));
System.out.println("生产者发送消息:"+ JSON.toJSONString(value));
//设置消息延迟级别,我这里设置5,对应就是延时一分钟
// "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h"
msg.setDelayTimeLevel(1);
msg.setDelayTimeLevel(level);
// 发送消息到一个Broker
result = producer.send(msg);
// 通过sendResult返回消息是否成功送达
@@ -1,8 +1,9 @@
package com.ruoyi.rocketmq.service;
import java.util.List;
import com.ruoyi.rocketmq.domain.InitialSwitchFanInfo;
import java.util.List;
/**
* 风扇信息Service接口
*
@@ -58,4 +59,11 @@ public interface IInitialSwitchFanInfoService
* @return 结果
*/
public int deleteInitialSwitchFanInfoById(Long id);
/**
* 批量新增风扇信息
* @param list
* @return
*/
public int insertBatchInitialSwitchFanInfo(List<InitialSwitchFanInfo> list);
}
@@ -1,8 +1,9 @@
package com.ruoyi.rocketmq.service;
import java.util.List;
import com.ruoyi.rocketmq.domain.InitialSwitchMpuInfo;
import java.util.List;
/**
* MPU信息Service接口
*
@@ -58,4 +59,11 @@ public interface IInitialSwitchMpuInfoService
* @return 结果
*/
public int deleteInitialSwitchMpuInfoById(Long id);
/**
* 批量新增mpu信息
* @param list
* @return
*/
public int insertBatchInitialSwitchMpuInfo(List<InitialSwitchMpuInfo> list);
}
@@ -1,8 +1,9 @@
package com.ruoyi.rocketmq.service;
import java.util.List;
import com.ruoyi.rocketmq.domain.InitialSwitchOpticalModule;
import java.util.List;
/**
* 光模块信息Service接口
*
@@ -58,4 +59,10 @@ public interface IInitialSwitchOpticalModuleService
* @return 结果
*/
public int deleteInitialSwitchOpticalModuleById(Long id);
/**
* 批量插入光模块信息
* @param list
*/
int batchInitialSwitchOpticalModule(List<InitialSwitchOpticalModule> list);
}
@@ -1,8 +1,9 @@
package com.ruoyi.rocketmq.service;
import java.util.List;
import com.ruoyi.rocketmq.domain.InitialSwitchPowerSupply;
import java.util.List;
/**
* 电源信息Service接口
*
@@ -58,4 +59,11 @@ public interface IInitialSwitchPowerSupplyService
* @return 结果
*/
public int deleteInitialSwitchPowerSupplyById(Long id);
/**
* 批量新增电源信息
* @param list
* @return
*/
public int insertBatchInitialSwitchPowerSupply(List<InitialSwitchPowerSupply> list);
}
@@ -93,4 +93,13 @@ public class InitialSwitchFanInfoServiceImpl implements IInitialSwitchFanInfoSer
{
return initialSwitchFanInfoMapper.deleteInitialSwitchFanInfoById(id);
}
/**
* 批量新增风扇信息
* @param list
* @return
*/
@Override
public int insertBatchInitialSwitchFanInfo(List<InitialSwitchFanInfo> list) {
return initialSwitchFanInfoMapper.insertBatchInitialSwitchFanInfo(list);
}
}
@@ -93,4 +93,14 @@ public class InitialSwitchMpuInfoServiceImpl implements IInitialSwitchMpuInfoSer
{
return initialSwitchMpuInfoMapper.deleteInitialSwitchMpuInfoById(id);
}
/**
* 批量新增mpu信息
* @param list
* @return
*/
@Override
public int insertBatchInitialSwitchMpuInfo(List<InitialSwitchMpuInfo> list) {
return initialSwitchMpuInfoMapper.insertBatchInitialSwitchMpuInfo(list);
}
}
@@ -93,4 +93,13 @@ public class InitialSwitchOpticalModuleServiceImpl implements IInitialSwitchOpti
{
return initialSwitchOpticalModuleMapper.deleteInitialSwitchOpticalModuleById(id);
}
/**
* 批量新增光模块信息
* @param list
*/
@Override
public int batchInitialSwitchOpticalModule(List<InitialSwitchOpticalModule> list) {
return initialSwitchOpticalModuleMapper.batchInitialSwitchOpticalModule(list);
}
}
@@ -93,4 +93,13 @@ public class InitialSwitchPowerSupplyServiceImpl implements IInitialSwitchPowerS
{
return initialSwitchPowerSupplyMapper.deleteInitialSwitchPowerSupplyById(id);
}
/**
* 批量新增电源信息
* @param list
* @return
*/
@Override
public int insertBatchInitialSwitchPowerSupply(List<InitialSwitchPowerSupply> list){
return initialSwitchPowerSupplyMapper.insertBatchInitialSwitchPowerSupply(list);
}
}
@@ -10,7 +10,6 @@ import com.ruoyi.rocketmq.mapper.RmAgentManagementMapper;
import com.ruoyi.rocketmq.model.ProducerMode;
import com.ruoyi.rocketmq.producer.MessageProducer;
import com.ruoyi.rocketmq.service.IRmAgentManagementService;
import com.ruoyi.system.api.RemoteRevenueConfigService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
@@ -35,13 +34,10 @@ public class RmAgentManagementServiceImpl implements IRmAgentManagementService
@Value("${fileDictory.filePath}")
private String filePath;
// private static String COMMAND = "/usr/local/tongran/deploy.sh restart";
// private static String COMMAND = "nohup /usr/local/tongran/deploy.sh restart > /dev/null 2>&1 &";
private static String COMMAND = "nohup /data/saas-income/houduan/tr-client-test/deploy.sh restart > /dev/null 2>&1 &";
private static String COMMAND2 = "cat /data/saas-income/houduan/tr-client-test/deploy.sh";
@Autowired
private ProducerMode producerMode;
@Autowired
private RemoteRevenueConfigService remoteRevenueConfigService;
/**
* 查询Agent管理
@@ -164,10 +160,9 @@ public class RmAgentManagementServiceImpl implements IRmAgentManagementService
// 构建更新策略
AgentUpdateVo agentUpdateVo = new AgentUpdateVo();
agentUpdateVo.setFileUrl(rmAgentManagement.getFileUrl());
agentUpdateVo.setFilePath("/data/saas-income/houduan/tr-client-test/temp");
agentUpdateVo.setFilePath(filePath);
List<String> commandList = new ArrayList<>();
commandList.add(COMMAND);
commandList.add(COMMAND2);
agentUpdateVo.setCommands(JSONObject.toJSONString(commandList));
agentUpdateVo.setMethod(rmAgentManagement.getMethod());
if(rmAgentManagement.getMethod() == 1){
@@ -300,7 +300,7 @@ public class RmMonitorPolicyServiceImpl implements IRmMonitorPolicyService
sendConfigurationsToDevices(devices, uniqueList);
}
// 更新策略状态为已下发
RmMonitorPolicy policyUpdate = new RmMonitorPolicy();
// RmMonitorPolicy policyUpdate = new RmMonitorPolicy();
// policyUpdate.setId(id);
// policyUpdate.setStatus("1");
// rmMonitorPolicyMapper.updateRmMonitorPolicy(policyUpdate);
@@ -467,28 +467,6 @@ public class RmMonitorPolicyServiceImpl implements IRmMonitorPolicyService
return switchOidVo;
}
// 将Map转换为字符串格式:{"oid1":"metric1","oid2":"metric2"}
private String mapToString(Map<String, String> map) {
if (map == null || map.isEmpty()) {
return "{}";
}
StringBuilder sb = new StringBuilder("{");
boolean first = true;
for (Map.Entry<String, String> entry : map.entrySet()) {
if (!first) {
sb.append(",");
}
sb.append("\"").append(entry.getKey()).append("\":\"").append(entry.getValue()).append("\"");
first = false;
}
sb.append("}");
return sb.toString();
}
/**
* 发送配置到设备
*/
@@ -498,19 +476,32 @@ public class RmMonitorPolicyServiceImpl implements IRmMonitorPolicyService
map.put("collects", collectVos);
map.put("timestamp", Instant.now().getEpochSecond());
String configJson = JSONObject.toJSONString(map);
Map stopMap = new HashMap();
map.put("timestamp", Instant.now().getEpochSecond());
for (RmResourceRegistrationRemote device : devices) {
try {
DeviceMessage message = new DeviceMessage();
message.setClientId(device.getHardwareSn());
message.setData(configJson);
message.setDataType(MsgEnum.开启系统采集.getValue());
DeviceMessage stopMsg = new DeviceMessage();
stopMsg.setClientId(device.getHardwareSn());
stopMsg.setData(JSONObject.toJSONString(stopMap));
stopMsg.setDataType(MsgEnum.关闭系统采集.getValue());
messageProducer.sendAsyncProducerMessage(
producerMode.getAgentTopic(),
"",
"",
JSONObject.toJSONString(message)
JSONObject.toJSONString(stopMsg)
);
DeviceMessage message = new DeviceMessage();
message.setClientId(device.getHardwareSn());
message.setData(configJson);
message.setDataType(MsgEnum.开启系统采集.getValue());
// 延迟1s防止同时读取关闭和开启
messageProducer.sendDelayMessage(
producerMode.getAgentTopic(),
"",
"",
JSONObject.toJSONString(message),1
);
} catch (Exception e) {
log.error("发送设备配置失败,deviceId: {}", device.getHardwareSn(), e);
@@ -525,6 +516,8 @@ public class RmMonitorPolicyServiceImpl implements IRmMonitorPolicyService
Map map = new HashMap();
map.put("collects", collectVos);
map.put("timestamp", Instant.now().getEpochSecond());
Map stopCollectMap = new HashMap();
map.put("timestamp", Instant.now().getEpochSecond());
String configJson = JSONObject.toJSONString(map);
for (RmResourceRegistrationRemote device : devices) {
@@ -553,17 +546,29 @@ public class RmMonitorPolicyServiceImpl implements IRmMonitorPolicyService
"",
JSONObject.toJSONString(registerMessage)
);
// 交换机采集
DeviceMessage message = new DeviceMessage();
message.setClientId(device.getHardwareSn());
message.setData(configJson);
message.setDataType(MsgEnum.开启交换机采集.getValue());
// 先停止采集,再开启采集,防止采集其他策略的信息
DeviceMessage stopCollectMsg = new DeviceMessage();
stopCollectMsg.setClientId(device.getHardwareSn());
stopCollectMsg.setData(JSONObject.toJSONString(stopCollectMap));
stopCollectMsg.setDataType(MsgEnum.关闭交换机采集.getValue());
// 待注册后下发策略,1秒延迟
messageProducer.sendDelayMessage(
producerMode.getAgentTopic(),
"",
"",
JSONObject.toJSONString(message)
JSONObject.toJSONString(stopCollectMsg),1
);
// 交换机采集
DeviceMessage message = new DeviceMessage();
message.setClientId(device.getHardwareSn());
message.setData(configJson);
message.setDataType(MsgEnum.开启交换机采集.getValue());
// 待注册后下发策略,5秒延迟
messageProducer.sendDelayMessage(
producerMode.getAgentTopic(),
"",
"",
JSONObject.toJSONString(message),2
);
} catch (Exception e) {
log.error("发送设备配置失败,deviceId: {}", device.getHardwareSn(), e);