diff --git a/tongran-common/tongran-common-core/src/main/java/com/tongran/common/core/utils/NetworkNameUtil.java b/tongran-common/tongran-common-core/src/main/java/com/tongran/common/core/utils/NetworkNameUtil.java index 35549cb..5eef0e6 100644 --- a/tongran-common/tongran-common-core/src/main/java/com/tongran/common/core/utils/NetworkNameUtil.java +++ b/tongran-common/tongran-common-core/src/main/java/com/tongran/common/core/utils/NetworkNameUtil.java @@ -13,8 +13,8 @@ public class NetworkNameUtil { return false; } - // 匹配模式:冒号后跟数字 或 点后跟数字 - String pattern = ".*[:.]\\d+$"; + // 匹配模式:冒号后跟数字 或 点 或 -后跟数字 + String pattern = ".*[:.-]\\d+$"; return interfaceName.matches(pattern); } } diff --git a/tongran-rocketmq/pom.xml b/tongran-rocketmq/pom.xml index bbfff4b..e753d87 100644 --- a/tongran-rocketmq/pom.xml +++ b/tongran-rocketmq/pom.xml @@ -93,6 +93,23 @@ snmp4j 2.8.9 + + + org.apache.iotdb + iotdb-jdbc + 1.3.2 + + + + org.apache.commons + commons-pool2 + 2.11.1 + + + org.apache.iotdb + iotdb-session + 1.3.2 + diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/config/SessionPoolConfig.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/config/SessionPoolConfig.java new file mode 100644 index 0000000..0b9004c --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/config/SessionPoolConfig.java @@ -0,0 +1,84 @@ +//package com.tongran.rocketmq.config; +// +//import lombok.extern.slf4j.Slf4j; +//import org.apache.iotdb.session.pool.SessionPool; +//import org.springframework.beans.factory.annotation.Value; +//import org.springframework.context.annotation.Bean; +//import org.springframework.context.annotation.Configuration; +// +//@Configuration +//@Slf4j +//public class SessionPoolConfig { +// +// +// /** +// * IoTDB连接配置 +// */ +// @Value("${iotdb.host:172.16.15.51}") +// private String host; +// +// @Value("${iotdb.port:6667}") +// private Integer port; +// +// @Value("${iotdb.username:root}") +// private String username; +// +// @Value("${iotdb.password:root123}") +// private String password; +// +// @Value("${iotdb.pool.maxSize:10}") +// private int maxSize; +// +// @Value("${iotdb.pool.enableCompression:false}") +// private boolean enableCompression; +// +// /** +// * 创建IoTDB SessionPool连接池 +// */ +// @Bean(destroyMethod = "close") +// public SessionPool sessionPool() { +// try { +// log.info("开始初始化IoTDB SessionPool连接池..."); +// log.info("IoTDB连接参数 - Host: {}, Port: {}, Username: {}, PoolSize: {}", +// host, port, username, maxSize); +// +// // 构建SessionPool +// SessionPool sessionPool = new SessionPool.Builder() +// .host(host) +// .port(port) +// .user(username) +// .password(password) +// .maxSize(maxSize) +// .enableCompression(enableCompression) +// .build(); +// +// // 测试连接 +// if (testConnection(sessionPool)) { +// log.info("IoTDB SessionPool连接池初始化成功!"); +// return sessionPool; +// } else { +// log.error("IoTDB连接测试失败,请检查配置和IoTDB服务状态"); +// throw new RuntimeException("IoTDB连接测试失败"); +// } +// } catch (Exception e) { +// log.error("初始化IoTDB SessionPool失败: ", e); +// throw new RuntimeException("初始化IoTDB SessionPool失败", e); +// } +// } +// +// /** +// * 测试IoTDB连接 +// */ +// private boolean testConnection(SessionPool sessionPool) { +// try { +// // 执行一个简单的查询来测试连接 +// String testSql = "show version"; +// sessionPool.executeQueryStatement(testSql); +// log.info("IoTDB连接测试成功"); +// return true; +// } catch (Exception e) { +// log.error("IoTDB连接测试异常: ", e); +// return false; +// } +// } +//} \ No newline at end of file diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/IotDbData.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/IotDbData.java new file mode 100644 index 0000000..cda0317 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/IotDbData.java @@ -0,0 +1,22 @@ +package com.tongran.rocketmq.domain; + +import lombok.Data; +import java.util.HashMap; +import java.util.Map; + +/** + * IoTDB数据包装类 + */ +@Data +public class IotDbData { + private String dataType; // 数据类型: server, switch, router + private String deviceId; // 设备标识 + private String tag; // 标签: 网卡名、端口名等 + private long timestamp; // 时间戳 + private Map dataMap = new HashMap<>(); // 数据字段 + + public IotDbData addField(String key, Object value) { + dataMap.put(key, value); + return this; + } +} \ No newline at end of file diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/vo/EpsMethodChangeRecordVo.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/vo/EpsMethodChangeRecordVo.java new file mode 100644 index 0000000..18fdf2d --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/vo/EpsMethodChangeRecordVo.java @@ -0,0 +1,23 @@ +package com.tongran.rocketmq.domain.vo; + +import com.tongran.common.core.web.domain.BaseEntity; +import lombok.Data; + +@Data +public class EpsMethodChangeRecordVo extends BaseEntity { + + + /** 修改内容 */ + private String changeContent; + + /** 创建人 */ + private String creatBy; + + /** 流量网口 */ + private String trafficPort; + + /** 业务代码(12位) */ + private String businessCode; + /** 客户端id */ + private String clientId; +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/MessageHandler.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/MessageHandler.java index fda46b3..676b026 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/MessageHandler.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/MessageHandler.java @@ -6,6 +6,7 @@ import com.tongran.common.core.domain.R; import com.tongran.common.core.enums.MsgEnum; import com.tongran.common.core.utils.DateUtils; import com.tongran.common.core.utils.StringUtils; +import com.tongran.common.security.utils.SecurityUtils; import com.tongran.rocketmq.domain.*; import com.tongran.rocketmq.domain.vo.*; import com.tongran.rocketmq.enums.AlarmTypeEnum; @@ -862,6 +863,11 @@ public class MessageHandler { redisTemplate.delete(statusKey); log.info("客户端ID: {} 已设置告警并清理心跳记录", clientId); + // frpc状态改为未运行 + RmFrpcConfigManage updateFrp = new RmFrpcConfigManage(); + updateFrp.setClientId(clientId); + updateFrp.setFrpcStatus("0"); + rmFrpcConfigManageService.updateRmFrpcConfigManage(updateFrp); // 修改资源状态 updateResourceStatus(clientId, "0"); } @@ -1343,6 +1349,24 @@ public class MessageHandler { updateData.setInterfaceType(networkInfo.getType()); needUpdate = true; } + if (!StringUtils.equals(networkInfo.getName(), oldInterfaceMsg.getInterfaceName())) { + if(networkInfo.getName() != null){ + // 添加业务变更记录 + EpsMethodChangeRecordVo recordAddData = new EpsMethodChangeRecordVo(); + recordAddData.setClientId(clientId); + recordAddData.setTrafficPort(networkInfo.getName()); + recordAddData.setUpdateTime(DateUtils.getNowDate()); + recordAddData.setCreateTime(DateUtils.getNowDate()); + recordAddData.setUpdateBy(SecurityUtils.getUsername()); + recordAddData.setCreatBy(SecurityUtils.getUsername()); + StringBuilder content = new StringBuilder(); + content.append("流量网口设置为").append(networkInfo.getName()); + recordAddData.setChangeContent(content.toString()); + rmNetworkInterfaceService.addTrafficPortChangeRecord(recordAddData); + updateData.setInterfaceName(networkInfo.getName()); + needUpdate = true; + } + } // 只有有字段变化时才执行更新 if (needUpdate) { diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/RmNetworkInterfaceMapper.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/RmNetworkInterfaceMapper.java index f68d2b8..33e8965 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/RmNetworkInterfaceMapper.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/RmNetworkInterfaceMapper.java @@ -1,6 +1,7 @@ package com.tongran.rocketmq.mapper; import com.tongran.rocketmq.domain.RmNetworkInterface; +import com.tongran.rocketmq.domain.vo.EpsMethodChangeRecordVo; import java.util.List; @@ -63,4 +64,6 @@ public interface RmNetworkInterfaceMapper int updateRmNetworkInterfaceByMac(RmNetworkInterface rmNetworkInterface); int updateNetMsgByMac(RmNetworkInterface rmNetworkInterface); + + int addTrafficPortChangeRecord(EpsMethodChangeRecordVo recordAddData); } diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IRmNetworkInterfaceService.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IRmNetworkInterfaceService.java index 18bce3c..97e7812 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IRmNetworkInterfaceService.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IRmNetworkInterfaceService.java @@ -1,6 +1,7 @@ package com.tongran.rocketmq.service; import com.tongran.rocketmq.domain.RmNetworkInterface; +import com.tongran.rocketmq.domain.vo.EpsMethodChangeRecordVo; import java.util.List; @@ -87,4 +88,6 @@ public interface IRmNetworkInterfaceService * @param clientId */ void issueNetName(String clientId); + + int addTrafficPortChangeRecord(EpsMethodChangeRecordVo recordAddData); } diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/RmNetworkInterfaceServiceImpl.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/RmNetworkInterfaceServiceImpl.java index ba2ddf1..f883872 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/RmNetworkInterfaceServiceImpl.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/RmNetworkInterfaceServiceImpl.java @@ -6,6 +6,7 @@ import com.tongran.common.core.utils.DateUtils; import com.tongran.rocketmq.domain.DeviceMessage; import com.tongran.rocketmq.domain.RmNetworkInterface; import com.tongran.rocketmq.domain.RmNetworkInterfaceChild; +import com.tongran.rocketmq.domain.vo.EpsMethodChangeRecordVo; import com.tongran.rocketmq.domain.vo.PolicyTypeVo; import com.tongran.rocketmq.mapper.RmNetworkInterfaceChildMapper; import com.tongran.rocketmq.mapper.RmNetworkInterfaceMapper; @@ -245,4 +246,9 @@ public class RmNetworkInterfaceServiceImpl implements IRmNetworkInterfaceService JSONObject.toJSONString(message) ); } + + @Override + public int addTrafficPortChangeRecord(EpsMethodChangeRecordVo recordAddData) { + return rmNetworkInterfaceMapper.addTrafficPortChangeRecord(recordAddData); + } } diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/snmp/service/ProcessSwitchCollectDataService.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/snmp/service/ProcessSwitchCollectDataService.java index db94af7..6e6dd34 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/snmp/service/ProcessSwitchCollectDataService.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/snmp/service/ProcessSwitchCollectDataService.java @@ -42,6 +42,8 @@ public class ProcessSwitchCollectDataService { private IInitialSwitchOpticalModuleService initialSwitchOpticalModuleService; @Autowired private IInitialSwitchOtherCollectDataService insertInitialSwitchOtherInfo; +// @Autowired +// private IotDbUtils iotDbUtils; // 在类中添加 private static final ObjectMapper objectMapper = new ObjectMapper(); /** @@ -198,6 +200,7 @@ public class ProcessSwitchCollectDataService { // 3. 计算速度 switchInfos.forEach(switchInfo -> { +// Map filedMap = new HashMap(); switchInfo.setClientId(clientId); switchInfo.setCreateTime(createTime); switchInfo.setSwitchIp(switchDataVo.getSwitchIp()); @@ -209,6 +212,7 @@ public class ProcessSwitchCollectDataService { // 字节转为bit BigDecimal inDiffBit = inDiff.multiply(new BigDecimal(8)).setScale(0, RoundingMode.HALF_UP); switchInfo.setInSpeed(inDiffBit.divide(divisor, 0, RoundingMode.HALF_UP)); +// filedMap.put("in", inDiffBit.divide(divisor, 0, RoundingMode.HALF_UP)); } // 计算outSpeed @@ -217,8 +221,12 @@ public class ProcessSwitchCollectDataService { // 字节转为bit BigDecimal outDiffBit = outDiff.multiply(new BigDecimal(8)).setScale(0, RoundingMode.HALF_UP); switchInfo.setOutSpeed(outDiffBit.divide(divisor, 0, RoundingMode.HALF_UP)); +// filedMap.put("out", outDiffBit.divide(divisor, 0, RoundingMode.HALF_UP)); } } +// iotDbUtils.writeDataWithTag("switch", clientId, +// switchInfo.getName(), millis, filedMap +// ); }); // 清空临时表对应switch信息 initialSwitchInfoTempService.truncateSwitchInfoTemp(clientId); diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/utils/IotDbUtils.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/utils/IotDbUtils.java new file mode 100644 index 0000000..588eb7f --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/utils/IotDbUtils.java @@ -0,0 +1,137 @@ +//package com.tongran.rocketmq.utils; +// +//import com.tongran.rocketmq.domain.IotDbData; +//import lombok.RequiredArgsConstructor; +//import lombok.extern.slf4j.Slf4j; +//import org.apache.iotdb.session.pool.SessionPool; +//import org.springframework.stereotype.Component; +// +//import java.util.List; +//import java.util.Map; +// +//@Component +//@Slf4j +//@RequiredArgsConstructor +//public class IotDbUtils { +// +// private final SessionPool sessionPool; +// +// /** +// * 通用写入方法 +// * @param dataType 数据类型: server, switch, router等 +// * @param deviceId 设备ID: clientId, switchIp等 +// * @param timestamp 时间戳 +// * @param dataMap 数据Map: key=字段名, value=字段值 +// */ +// public void writeData(String dataType, String deviceId, long timestamp, Map dataMap) { +// if (dataMap == null || dataMap.isEmpty()) return; +// +// try { +// String sql = buildInsertSql(dataType, deviceId, timestamp, dataMap); +// sessionPool.executeNonQueryStatement(sql); +// } catch (Exception e) { +// log.error("IoTDB写入失败: type={}, device={}", dataType, deviceId, e); +// } +// } +// +// /** +// * 写入带标签的数据(如网卡、端口) +// */ +// public void writeDataWithTag(String dataType, String deviceId, String tag, +// long timestamp, Map dataMap) { +// if (dataMap == null || dataMap.isEmpty()) return; +// +// try { +// String sql = buildInsertSqlWithTag(dataType, deviceId, tag, timestamp, dataMap); +// sessionPool.executeNonQueryStatement(sql); +// } catch (Exception e) { +// log.error("IoTDB写入失败: type={}, device={}, tag={}", dataType, deviceId, tag, e); +// } +// } +// +// /** +// * 批量写入 +// */ +// public void batchWrite(List dataList) { +// if (dataList == null || dataList.isEmpty()) return; +// +// for (IotDbData data : dataList) { +// writeData(data.getDataType(), data.getDeviceId(), +// data.getTimestamp(), data.getDataMap()); +// } +// } +// +// private String buildInsertSql(String dataType, String deviceId, +// long timestamp, Map dataMap) { +// String devicePath = buildDevicePath(dataType, deviceId, null); +// return buildSql(devicePath, timestamp, dataMap); +// } +// +// private String buildInsertSqlWithTag(String dataType, String deviceId, String tag, +// long timestamp, Map dataMap) { +// String devicePath = buildDevicePath(dataType, deviceId, tag); +// return buildSql(devicePath, timestamp, dataMap); +// } +// +// private String buildDevicePath(String dataType, String deviceId, String tag) { +// // 基础路径: root.{dataType}.{deviceId} +// StringBuilder path = new StringBuilder("root.") +// .append(safePath(dataType)).append(".") +// .append(safePath(deviceId)); +// +// // 如果有标签(如网卡名、端口名) +// if (tag != null && !tag.trim().isEmpty()) { +// path.append(".").append(safePath(tag)); +// } +// +// return path.toString(); +// } +// +// private String buildSql(String devicePath, long timestamp, Map dataMap) { +// StringBuilder sql = new StringBuilder("INSERT INTO ") +// .append(devicePath) +// .append("(timestamp"); +// +// StringBuilder values = new StringBuilder("VALUES(").append(timestamp); +// +// for (Map.Entry entry : dataMap.entrySet()) { +// if (entry.getValue() != null) { +// sql.append(", ").append(safeField(entry.getKey())); +// values.append(", ").append(formatValue(entry.getValue())); +// } +// } +// +// sql.append(") ").append(values).append(")"); +// return sql.toString(); +// } +// +// /** +// * 路径安全处理 +// */ +// private String safePath(String input) { +// if (input == null) return "unknown"; +// return input.replace(".", "_") +// .replace(":", "_") +// .replace("-", "_") +// .replaceAll("[^a-zA-Z0-9_]", "_"); +// } +// +// /** +// * 字段名安全处理 +// */ +// private String safeField(String field) { +// if (field == null) return "unknown"; +// return field.replaceAll("[^a-zA-Z0-9_]", "_"); +// } +// +// /** +// * 值格式化 +// */ +// private String formatValue(Object value) { +// if (value == null) return "null"; +// if (value instanceof String) { +// return "'" + value.toString().replace("'", "''") + "'"; +// } +// return value.toString(); +// } +//} \ No newline at end of file diff --git a/tongran-rocketmq/src/main/resources/mapper/rocketmq/RmFrpcConfigManageMapper.xml b/tongran-rocketmq/src/main/resources/mapper/rocketmq/RmFrpcConfigManageMapper.xml index a233fd7..cd87e9d 100644 --- a/tongran-rocketmq/src/main/resources/mapper/rocketmq/RmFrpcConfigManageMapper.xml +++ b/tongran-rocketmq/src/main/resources/mapper/rocketmq/RmFrpcConfigManageMapper.xml @@ -104,7 +104,6 @@ update rm_frpc_config_manage - client_id = #{clientId}, frpc_server_addr = #{frpcServerAddr}, frpc_server_port = #{frpcServerPort}, frpc_remote_port = #{frpcRemotePort}, @@ -120,7 +119,19 @@ frpc_status = #{frpcStatus}, frpc_type = #{frpcType}, - where id = #{id} + + + + and id = #{id} + + + and client_id = #{clientId} + + + and 1=0 + + + diff --git a/tongran-rocketmq/src/main/resources/mapper/rocketmq/RmNetworkInterfaceMapper.xml b/tongran-rocketmq/src/main/resources/mapper/rocketmq/RmNetworkInterfaceMapper.xml index e7ef010..1d4a381 100644 --- a/tongran-rocketmq/src/main/resources/mapper/rocketmq/RmNetworkInterfaceMapper.xml +++ b/tongran-rocketmq/src/main/resources/mapper/rocketmq/RmNetworkInterfaceMapper.xml @@ -176,4 +176,25 @@ update_time = #{updateTime} where mac_address = #{macAddress} and client_id = #{clientId} + + insert into eps_method_change_record + + change_content, + create_time, + creat_by, + update_time, + update_by, + traffic_port, + client_id, + + + #{changeContent}, + #{createTime}, + #{creatBy}, + #{updateTime}, + #{updateBy}, + #{trafficPort}, + #{clientId}, + + \ No newline at end of file