From 2593c587416546b11c638b2b54f986d84e8c185e Mon Sep 17 00:00:00 2001 From: gaoyutao Date: Sun, 4 Jan 2026 18:10:13 +0800 Subject: [PATCH] =?UTF-8?q?=E6=9C=8D=E5=8A=A1=E5=99=A8=E7=B3=BB=E7=BB=9F?= =?UTF-8?q?=E5=85=B6=E4=BB=96=E4=BF=A1=E6=81=AF=E6=95=B0=E6=8D=AE=E8=A1=A8?= =?UTF-8?q?=E5=88=86=E8=A1=A8=E3=80=82=20=E4=BA=A4=E6=8D=A2=E6=9C=BA?= =?UTF-8?q?=E5=85=89=E6=A8=A1=E5=9D=97=E4=BF=A1=E6=81=AF=E6=95=B0=E6=8D=AE?= =?UTF-8?q?=E8=A1=A8=E5=88=86=E8=A1=A8=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../mapper/EpsInitialTrafficDataMapper.java | 3 + .../EpsInitialTrafficDataServiceImpl.java | 15 +++ .../RmResourceRegistrationServiceImpl.java | 6 +- .../system/EpsInitialTrafficDataMapper.xml | 33 ++++- .../rocketmq/domain/AllMoudleName.java | 69 +++++++++++ .../domain/InitialSwitchOpticalModule.java | 5 + .../domain/InitialSystemOtherCollectData.java | 2 + .../handler/DeviceMessageHandler.java | 10 +- .../rocketmq/mapper/AllMoudleNameMapper.java | 64 ++++++++++ .../InitialSwitchOpticalModuleMapper.java | 4 +- .../service/IAllMoudleNameService.java | 65 ++++++++++ .../IInitialSwitchOpticalModuleService.java | 3 +- .../impl/AllMoudleNameServiceImpl.java | 116 ++++++++++++++++++ ...InitialSwitchOpticalModuleServiceImpl.java | 85 +++++++++++-- ...tialSystemOtherCollectDataServiceImpl.java | 46 +++++-- .../ProcessSwitchCollectDataService.java | 10 +- .../mapper/rocketmq/AllMoudleNameMapper.xml | 82 +++++++++++++ .../InitialSwitchOpticalModuleMapper.xml | 18 +-- .../InitialSystemOtherCollectDataMapper.xml | 10 +- 19 files changed, 597 insertions(+), 49 deletions(-) create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/AllMoudleName.java create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/AllMoudleNameMapper.java create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IAllMoudleNameService.java create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/AllMoudleNameServiceImpl.java create mode 100644 tongran-rocketmq/src/main/resources/mapper/rocketmq/AllMoudleNameMapper.xml diff --git a/tongran-modules/tongran-system/src/main/java/com/tongran/system/mapper/EpsInitialTrafficDataMapper.java b/tongran-modules/tongran-system/src/main/java/com/tongran/system/mapper/EpsInitialTrafficDataMapper.java index c2dfaed..70c14bb 100644 --- a/tongran-modules/tongran-system/src/main/java/com/tongran/system/mapper/EpsInitialTrafficDataMapper.java +++ b/tongran-modules/tongran-system/src/main/java/com/tongran/system/mapper/EpsInitialTrafficDataMapper.java @@ -69,4 +69,7 @@ public interface EpsInitialTrafficDataMapper { List getTrafficListByClientIds(EpsInitialTrafficData condition); + void createOtherMsgTable(String tableName); + + void createSwitchOpMdTable(String tableName); } diff --git a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/EpsInitialTrafficDataServiceImpl.java b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/EpsInitialTrafficDataServiceImpl.java index 6628239..f621537 100644 --- a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/EpsInitialTrafficDataServiceImpl.java +++ b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/EpsInitialTrafficDataServiceImpl.java @@ -70,6 +70,10 @@ public class EpsInitialTrafficDataServiceImpl implements EpsInitialTrafficDataSe createSwitchInfoTable(year, month); // 创建交换机流量业务表 createSwitchInfoDetailsTable(year, month); + // 创建服务器其他信息表 + createOtherMsgTable(year, month); + // 创建交换机光模块表 + createSwitchOpMdTable(year, month); } private void createTrafficDetailsTable(int year, int month) { @@ -100,6 +104,17 @@ public class EpsInitialTrafficDataServiceImpl implements EpsInitialTrafficDataSe }); } + private void createOtherMsgTable(int year, int month) { + createRangeTables(year, month, "initial_system_other_collect_data", (tableName) -> { + epsInitialTrafficDataMapper.createOtherMsgTable(tableName); + }); + } + private void createSwitchOpMdTable(int year, int month) { + createRangeTables(year, month, "initial_switch_optical_module", (tableName) -> { + epsInitialTrafficDataMapper.createSwitchOpMdTable(tableName); + }); + } + /** * 通用创建表方法 diff --git a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/RmResourceRegistrationServiceImpl.java b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/RmResourceRegistrationServiceImpl.java index a5d0ba2..f71543a 100644 --- a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/RmResourceRegistrationServiceImpl.java +++ b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/RmResourceRegistrationServiceImpl.java @@ -1159,9 +1159,9 @@ public class RmResourceRegistrationServiceImpl implements IRmResourceRegistratio String gateWay = mtmNetworks.get(0).getGateway(); routeMsg.setName(mtmNetworks.get(0).getInterfaceName()); // 使用拼接后的网络接口名称 routeMsg.setGateway(gateWay); // 使用第一个网络接口的网关,或者您可以根据需要调整 - if(gateWay != null){ - rspVo.setAddRoute(JSONObject.toJSONString(routeMsg)); - } +// if(gateWay != null){ +// rspVo.setAddRoute(JSONObject.toJSONString(routeMsg)); +// } } rspVo.setNetName(netName); diff --git a/tongran-modules/tongran-system/src/main/resources/mapper/system/EpsInitialTrafficDataMapper.xml b/tongran-modules/tongran-system/src/main/resources/mapper/system/EpsInitialTrafficDataMapper.xml index 7411df4..9b4e097 100644 --- a/tongran-modules/tongran-system/src/main/resources/mapper/system/EpsInitialTrafficDataMapper.xml +++ b/tongran-modules/tongran-system/src/main/resources/mapper/system/EpsInitialTrafficDataMapper.xml @@ -169,7 +169,38 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" INDEX idx_client_time (client_id, name, create_time) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='交换机监控详情表'; - + + CREATE TABLE IF NOT EXISTS ${tableName} ( + id BIGINT(20) NOT NULL AUTO_INCREMENT COMMENT '主键ID', + client_id VARCHAR(64) COMMENT '客户端ID', + collect_type VARCHAR(100) COMMENT '采集类型', + collect_value LONGTEXT COMMENT '采集值', + create_time DATETIME COMMENT '创建时间', + update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '修改时间', + create_by VARCHAR(64) COMMENT '创建人', + update_by VARCHAR(64) COMMENT '修改人', + PRIMARY KEY (id), + UNIQUE KEY uk_client_collect (client_id, collect_type, create_time) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='服务器系统其他信息采集数据表'; + + + CREATE TABLE IF NOT EXISTS ${tableName} ( + id BIGINT(20) NOT NULL AUTO_INCREMENT COMMENT '主键ID', + fiber_ent_index VARCHAR(64) COMMENT '光模块端口索引', + fiber_port_name VARCHAR(128) COMMENT '光模块端口名称', + hw_entity_optical_tx_low_threshold DECIMAL(10,2) COMMENT '光模块发送光衰阈值(dBm)', + hw_entity_optical_rx_low_threshold DECIMAL(10,2) COMMENT '光模块接收光衰阈值(dBm)', + hw_entity_optical_rx_power DECIMAL(10,2) COMMENT '光模块接收功率(dBm)', + hw_entity_optical_tx_power DECIMAL(10,2) COMMENT '光模块发送功率(dBm)', + client_id VARCHAR(64) COMMENT '客户端ID', + create_time DATETIME COMMENT '创建时间', + update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '修改时间', + create_by VARCHAR(64) COMMENT '创建人', + update_by VARCHAR(64) COMMENT '修改人', + PRIMARY KEY (id), + UNIQUE KEY uk_client_fiber_time (client_id, create_time, fiber_port_name) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci COMMENT='光模块信息表'; + INSERT INTO ${tableName} ( diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/AllMoudleName.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/AllMoudleName.java new file mode 100644 index 0000000..8995773 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/AllMoudleName.java @@ -0,0 +1,69 @@ +package com.tongran.rocketmq.domain; + +import org.apache.commons.lang3.builder.ToStringBuilder; +import org.apache.commons.lang3.builder.ToStringStyle; +import com.tongran.common.core.annotation.Excel; +import com.tongran.common.core.web.domain.BaseEntity; + +/** + * 光模块名称存储对象 all_moudle_name + * + * @author tongran + * @date 2026-01-04 + */ +public class AllMoudleName extends BaseEntity +{ + private static final long serialVersionUID = 1L; + + /** 自增主键 */ + private Long id; + + /** 客户端ID */ + @Excel(name = "客户端ID") + private String clientId; + + /** 光模块名称 */ + @Excel(name = "光模块名称") + private String name; + + public void setId(Long id) + { + this.id = id; + } + + public Long getId() + { + return id; + } + + public void setClientId(String clientId) + { + this.clientId = clientId; + } + + public String getClientId() + { + return clientId; + } + + public void setName(String name) + { + this.name = name; + } + + public String getName() + { + return name; + } + + @Override + public String toString() { + return new ToStringBuilder(this,ToStringStyle.MULTI_LINE_STYLE) + .append("id", getId()) + .append("clientId", getClientId()) + .append("name", getName()) + .append("createTime", getCreateTime()) + .append("updateTime", getUpdateTime()) + .toString(); + } +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialSwitchOpticalModule.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialSwitchOpticalModule.java index 7546d7b..f1d3ba1 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialSwitchOpticalModule.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialSwitchOpticalModule.java @@ -5,6 +5,7 @@ import com.tongran.common.core.web.domain.BaseEntity; import lombok.Data; import java.math.BigDecimal; +import java.util.List; /** * 光模块信息对象 initial_switch_optical_module @@ -51,4 +52,8 @@ public class InitialSwitchOpticalModule extends BaseEntity private String startTime; /** 结束时间 */ private String endTime; + /** 表名 */ + private String tableName; + + private List list; } diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialSystemOtherCollectData.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialSystemOtherCollectData.java index 198b611..ab256ba 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialSystemOtherCollectData.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialSystemOtherCollectData.java @@ -33,4 +33,6 @@ public class InitialSystemOtherCollectData extends BaseEntity private String startTime; /** 结束时间 */ private String endTime; + /** 表名 */ + private String tableName; } diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/DeviceMessageHandler.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/DeviceMessageHandler.java index c917000..a97a290 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/DeviceMessageHandler.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/DeviceMessageHandler.java @@ -370,15 +370,15 @@ public class DeviceMessageHandler { private void handleSwitchModuleMessage(CollectDataVo switchDataVo, String clientId){ List moduleList = SwitchJsonDataParser.parseJsonData(switchDataVo.getValue(), InitialSwitchOpticalModule.class); if (!moduleList.isEmpty()){ + // 时间戳转换 + long timestamp = switchDataVo.getTimestamp(); + long millis = timestamp * 1000; + Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒 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); + initialSwitchOpticalModuleService.batchInitialSwitchOpticalModule(moduleList, createTime); } } diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/AllMoudleNameMapper.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/AllMoudleNameMapper.java new file mode 100644 index 0000000..5452151 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/AllMoudleNameMapper.java @@ -0,0 +1,64 @@ +package com.tongran.rocketmq.mapper; + +import com.tongran.rocketmq.domain.AllMoudleName; + +import java.util.List; + +/** + * 光模块名称存储Mapper接口 + * + * @author tongran + * @date 2026-01-04 + */ +public interface AllMoudleNameMapper +{ + /** + * 查询光模块名称存储 + * + * @param id 光模块名称存储主键 + * @return 光模块名称存储 + */ + public AllMoudleName selectAllMoudleNameById(Long id); + + /** + * 查询光模块名称存储列表 + * + * @param allMoudleName 光模块名称存储 + * @return 光模块名称存储集合 + */ + public List selectAllMoudleNameList(AllMoudleName allMoudleName); + + /** + * 新增光模块名称存储 + * + * @param allMoudleName 光模块名称存储 + * @return 结果 + */ + public int insertAllMoudleName(AllMoudleName allMoudleName); + + /** + * 修改光模块名称存储 + * + * @param allMoudleName 光模块名称存储 + * @return 结果 + */ + public int updateAllMoudleName(AllMoudleName allMoudleName); + + /** + * 删除光模块名称存储 + * + * @param id 光模块名称存储主键 + * @return 结果 + */ + public int deleteAllMoudleNameById(Long id); + + /** + * 批量删除光模块名称存储 + * + * @param ids 需要删除的数据主键集合 + * @return 结果 + */ + public int deleteAllMoudleNameByIds(Long[] ids); + + void batchInsertMoudleName(List nameList); +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/InitialSwitchOpticalModuleMapper.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/InitialSwitchOpticalModuleMapper.java index 2da794b..eaeffc3 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/InitialSwitchOpticalModuleMapper.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/InitialSwitchOpticalModuleMapper.java @@ -63,9 +63,9 @@ public interface InitialSwitchOpticalModuleMapper /** * 批量新增光模块信息 - * @param list + * @param initialSwitchOpticalModule */ - int batchInitialSwitchOpticalModule(List list); + int batchInitialSwitchOpticalModule(InitialSwitchOpticalModule initialSwitchOpticalModule); /** * 光模块基础信息 diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IAllMoudleNameService.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IAllMoudleNameService.java new file mode 100644 index 0000000..ba8ac2a --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IAllMoudleNameService.java @@ -0,0 +1,65 @@ +package com.tongran.rocketmq.service; + +import com.tongran.rocketmq.domain.AllMoudleName; +import com.tongran.rocketmq.domain.InitialSwitchOpticalModule; + +import java.util.List; + +/** + * 光模块名称存储Service接口 + * + * @author tongran + * @date 2026-01-04 + */ +public interface IAllMoudleNameService +{ + /** + * 查询光模块名称存储 + * + * @param id 光模块名称存储主键 + * @return 光模块名称存储 + */ + public AllMoudleName selectAllMoudleNameById(Long id); + + /** + * 查询光模块名称存储列表 + * + * @param allMoudleName 光模块名称存储 + * @return 光模块名称存储集合 + */ + public List selectAllMoudleNameList(AllMoudleName allMoudleName); + + /** + * 新增光模块名称存储 + * + * @param allMoudleName 光模块名称存储 + * @return 结果 + */ + public int insertAllMoudleName(AllMoudleName allMoudleName); + + /** + * 修改光模块名称存储 + * + * @param allMoudleName 光模块名称存储 + * @return 结果 + */ + public int updateAllMoudleName(AllMoudleName allMoudleName); + + /** + * 批量删除光模块名称存储 + * + * @param ids 需要删除的光模块名称存储主键集合 + * @return 结果 + */ + public int deleteAllMoudleNameByIds(Long[] ids); + + /** + * 删除光模块名称存储信息 + * + * @param id 光模块名称存储主键 + * @return 结果 + */ + public int deleteAllMoudleNameById(Long id); + + void batchInsertMoudleName(List dataList); +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IInitialSwitchOpticalModuleService.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IInitialSwitchOpticalModuleService.java index d4682ec..b5a8c63 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IInitialSwitchOpticalModuleService.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IInitialSwitchOpticalModuleService.java @@ -2,6 +2,7 @@ package com.tongran.rocketmq.service; import com.tongran.rocketmq.domain.InitialSwitchOpticalModule; +import java.util.Date; import java.util.List; import java.util.Map; @@ -65,7 +66,7 @@ public interface IInitialSwitchOpticalModuleService * 批量插入光模块信息 * @param list */ - int batchInitialSwitchOpticalModule(List list); + int batchInitialSwitchOpticalModule(List list, Date createTime); /** * 光模块基础信息 diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/AllMoudleNameServiceImpl.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/AllMoudleNameServiceImpl.java new file mode 100644 index 0000000..2115488 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/AllMoudleNameServiceImpl.java @@ -0,0 +1,116 @@ +package com.tongran.rocketmq.service.impl; + +import com.tongran.common.core.utils.DateUtils; +import com.tongran.rocketmq.domain.AllMoudleName; +import com.tongran.rocketmq.domain.InitialSwitchOpticalModule; +import com.tongran.rocketmq.mapper.AllMoudleNameMapper; +import com.tongran.rocketmq.service.IAllMoudleNameService; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; + +import java.util.ArrayList; +import java.util.List; + +/** + * 光模块名称存储Service业务层处理 + * + * @author tongran + * @date 2026-01-04 + */ +@Service +public class AllMoudleNameServiceImpl implements IAllMoudleNameService +{ + @Autowired + private AllMoudleNameMapper allMoudleNameMapper; + + /** + * 查询光模块名称存储 + * + * @param id 光模块名称存储主键 + * @return 光模块名称存储 + */ + @Override + public AllMoudleName selectAllMoudleNameById(Long id) + { + return allMoudleNameMapper.selectAllMoudleNameById(id); + } + + /** + * 查询光模块名称存储列表 + * + * @param allMoudleName 光模块名称存储 + * @return 光模块名称存储 + */ + @Override + public List selectAllMoudleNameList(AllMoudleName allMoudleName) + { + return allMoudleNameMapper.selectAllMoudleNameList(allMoudleName); + } + + /** + * 新增光模块名称存储 + * + * @param allMoudleName 光模块名称存储 + * @return 结果 + */ + @Override + public int insertAllMoudleName(AllMoudleName allMoudleName) + { + allMoudleName.setCreateTime(DateUtils.getNowDate()); + return allMoudleNameMapper.insertAllMoudleName(allMoudleName); + } + + /** + * 修改光模块名称存储 + * + * @param allMoudleName 光模块名称存储 + * @return 结果 + */ + @Override + public int updateAllMoudleName(AllMoudleName allMoudleName) + { + allMoudleName.setUpdateTime(DateUtils.getNowDate()); + return allMoudleNameMapper.updateAllMoudleName(allMoudleName); + } + + /** + * 批量删除光模块名称存储 + * + * @param ids 需要删除的光模块名称存储主键 + * @return 结果 + */ + @Override + public int deleteAllMoudleNameByIds(Long[] ids) + { + return allMoudleNameMapper.deleteAllMoudleNameByIds(ids); + } + + /** + * 删除光模块名称存储信息 + * + * @param id 光模块名称存储主键 + * @return 结果 + */ + @Override + public int deleteAllMoudleNameById(Long id) + { + return allMoudleNameMapper.deleteAllMoudleNameById(id); + } + + @Override + public void batchInsertMoudleName(List dataList) { + if(dataList == null){ + dataList = new ArrayList<>(); + } + List nameList = new ArrayList<>(); + for (InitialSwitchOpticalModule initialSwitchOpticalModule : dataList) { + AllMoudleName allMoudleName = new AllMoudleName(); + allMoudleName.setClientId(initialSwitchOpticalModule.getClientId()); + allMoudleName.setName(initialSwitchOpticalModule.getFiberPortName()); + allMoudleName.setCreateTime(initialSwitchOpticalModule.getCreateTime()); + allMoudleName.setUpdateTime(DateUtils.getNowDate()); + nameList.add(allMoudleName); + } + allMoudleNameMapper.batchInsertMoudleName(nameList); + } +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialSwitchOpticalModuleServiceImpl.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialSwitchOpticalModuleServiceImpl.java index c931f3b..dd352af 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialSwitchOpticalModuleServiceImpl.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialSwitchOpticalModuleServiceImpl.java @@ -2,16 +2,21 @@ package com.tongran.rocketmq.service.impl; import com.tongran.common.core.utils.DateUtils; import com.tongran.common.core.utils.EchartsDataUtils; +import com.tongran.common.core.utils.TableSubUtil; import com.tongran.rocketmq.domain.InitialSwitchOpticalModule; import com.tongran.rocketmq.mapper.InitialSwitchOpticalModuleMapper; +import com.tongran.rocketmq.service.IAllMoudleNameService; import com.tongran.rocketmq.service.IInitialSwitchOpticalModuleService; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.BeanUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Isolation; +import org.springframework.transaction.annotation.Transactional; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.function.Function; +import java.util.stream.Collectors; /** * 光模块信息Service业务层处理 @@ -20,10 +25,15 @@ import java.util.function.Function; * @date 2025-09-22 */ @Service +@Slf4j public class InitialSwitchOpticalModuleServiceImpl implements IInitialSwitchOpticalModuleService { @Autowired private InitialSwitchOpticalModuleMapper initialSwitchOpticalModuleMapper; + @Autowired + private IAllMoudleNameService allMoudleNameService; + + private final static String TABLE_PREFIX = "initial_switch_optical_module"; /** * 查询光模块信息 @@ -98,13 +108,72 @@ public class InitialSwitchOpticalModuleServiceImpl implements IInitialSwitchOpti return initialSwitchOpticalModuleMapper.deleteInitialSwitchOpticalModuleById(id); } + /** + * 分表查询服务器其他信息 + * @param queryParam + * @return + */ + public List getSwitchOpticalModuleSharding(InitialSwitchOpticalModule queryParam) { + // 获取涉及的表名 + Set tableNames = TableSubUtil.getExistingTableNamesBetween(queryParam.getStartTime(), queryParam.getEndTime(), TABLE_PREFIX); + + // 并行查询各表 + return tableNames.parallelStream() + .flatMap(tableName -> { + InitialSwitchOpticalModule condition = new InitialSwitchOpticalModule(); + condition.setTableName(tableName); + condition.setClientId(queryParam.getClientId()); + condition.setFiberPortName(queryParam.getFiberPortName()); + condition.setStartTime(queryParam.getStartTime()); + condition.setEndTime(queryParam.getEndTime()); + return initialSwitchOpticalModuleMapper.selectInitialSwitchOpticalModuleList(condition).stream(); + }) + .collect(Collectors.toList()); + } /** * 批量新增光模块信息 * @param list */ @Override - public int batchInitialSwitchOpticalModule(List list) { - return initialSwitchOpticalModuleMapper.batchInitialSwitchOpticalModule(list); + @Transactional(rollbackFor = Exception.class, isolation = Isolation.READ_COMMITTED) + public int batchInitialSwitchOpticalModule(List list, Date createTime) { + if (list == null || list.isEmpty()) { + return 0; + } + // 按表名分组批量插入 + Map> groupedData = list.stream() + .map(data -> { + try { + InitialSwitchOpticalModule processed = new InitialSwitchOpticalModule(); + BeanUtils.copyProperties(data,processed); + if (data.getCreateTime() == null) { + data.setCreateTime(DateUtils.getNowDate()); + } + processed.setTableName(TableSubUtil.getTableName(createTime, TABLE_PREFIX)); + return processed; + } catch (Exception e){ + log.error("数据处理失败",e.getMessage()); + return null; + } + }).collect(Collectors.groupingBy( + InitialSwitchOpticalModule::getTableName, + LinkedHashMap::new, // 保持插入顺序 + Collectors.toList())); + + groupedData.forEach((tableName, dataList) -> { + try { + InitialSwitchOpticalModule data = new InitialSwitchOpticalModule(); + data.setTableName(tableName); + data.setList(dataList); + initialSwitchOpticalModuleMapper.batchInitialSwitchOpticalModule(data); + // 记录光模块名称 + allMoudleNameService.batchInsertMoudleName(dataList); + } catch (Exception e) { + log.error("表{}插入失败", tableName, e); + throw new RuntimeException("批量插入失败", e); + } + }); + return 1; } /** @@ -114,6 +183,8 @@ public class InitialSwitchOpticalModuleServiceImpl implements IInitialSwitchOpti */ @Override public List switchOpticalModuleMsg(InitialSwitchOpticalModule initialSwitchOpticalModule) { + String tableName = TableSubUtil.getTableName(DateUtils.getNowDate(), TABLE_PREFIX); + initialSwitchOpticalModule.setTableName(tableName); return initialSwitchOpticalModuleMapper.switchOpticalModuleMsg(initialSwitchOpticalModule); } @@ -134,7 +205,7 @@ public class InitialSwitchOpticalModuleServiceImpl implements IInitialSwitchOpti */ @Override public Map opticalModuleLowThreshold(InitialSwitchOpticalModule initialSwitchOpticalModule) { - List list = initialSwitchOpticalModuleMapper.selectInitialSwitchOpticalModuleList(initialSwitchOpticalModule); + List list = getSwitchOpticalModuleSharding(initialSwitchOpticalModule); Map> extractors = new LinkedHashMap<>(); extractors.put("TxLowThreshold", info -> info.getHwEntityOpticalTxLowThreshold()); extractors.put("RxLowThreshold", info -> info.getHwEntityOpticalRxLowThreshold()); @@ -148,7 +219,7 @@ public class InitialSwitchOpticalModuleServiceImpl implements IInitialSwitchOpti */ @Override public Map opticalModulePower(InitialSwitchOpticalModule initialSwitchOpticalModule) { - List list = initialSwitchOpticalModuleMapper.selectInitialSwitchOpticalModuleList(initialSwitchOpticalModule); + List list = getSwitchOpticalModuleSharding(initialSwitchOpticalModule); Map> extractors = new LinkedHashMap<>(); extractors.put("TxPower", info -> info.getHwEntityOpticalTxPower()); extractors.put("RxPower", info -> info.getHwEntityOpticalRxPower()); diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialSystemOtherCollectDataServiceImpl.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialSystemOtherCollectDataServiceImpl.java index 7ec1835..77f8c1b 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialSystemOtherCollectDataServiceImpl.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialSystemOtherCollectDataServiceImpl.java @@ -2,6 +2,7 @@ package com.tongran.rocketmq.service.impl; import com.tongran.common.core.utils.DateUtils; import com.tongran.common.core.utils.EchartsDataUtils; +import com.tongran.common.core.utils.TableSubUtil; import com.tongran.common.core.utils.UnitChangeUtil; import com.tongran.rocketmq.domain.InitialCpuInfo; import com.tongran.rocketmq.domain.InitialSystemOtherCollectData; @@ -30,6 +31,7 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO private InitialSystemOtherCollectDataMapper initialSystemOtherCollectDataMapper; @Autowired private InitialCpuInfoMapper initialCpuInfoMapper; + private final static String TABLE_PREFIX = "initial_system_other_collect_data"; /** * 查询交换机系统其他信息采集数据 @@ -64,6 +66,8 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO @Override public int insertInitialSystemOtherCollectData(InitialSystemOtherCollectData initialSystemOtherCollectData) { + String tableName = TableSubUtil.getTableName(initialSystemOtherCollectData.getCreateTime(), TABLE_PREFIX); + initialSystemOtherCollectData.setTableName(tableName); return initialSystemOtherCollectDataMapper.insertInitialSystemOtherCollectData(initialSystemOtherCollectData); } @@ -111,6 +115,8 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO */ @Override public Map getMonitorMsg(InitialSystemOtherCollectData initialSystemOtherCollectData) { + String tableName = TableSubUtil.getTableName(DateUtils.getNowDate(), TABLE_PREFIX); + initialSystemOtherCollectData.setTableName(tableName); Map map = initialSystemOtherCollectDataMapper.getMonitorMsg(initialSystemOtherCollectData); // 如果返回null,初始化一个空的Map if (map == null) { @@ -147,6 +153,28 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO map.put("user", cpuInfo.getUser()); return map; } + /** + * 分表查询服务器其他信息 + * @param queryParam + * @return + */ + public List getOtherMsgSharding(InitialSystemOtherCollectData queryParam) { + // 获取涉及的表名 + Set tableNames = TableSubUtil.getExistingTableNamesBetween(queryParam.getStartTime(), queryParam.getEndTime(), TABLE_PREFIX); + + // 并行查询各表 + return tableNames.parallelStream() + .flatMap(tableName -> { + InitialSystemOtherCollectData condition = new InitialSystemOtherCollectData(); + condition.setTableName(tableName); + condition.setClientId(queryParam.getClientId()); + condition.setCollectType(queryParam.getCollectType()); + condition.setStartTime(queryParam.getStartTime()); + condition.setEndTime(queryParam.getEndTime()); + return initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(condition).stream(); + }) + .collect(Collectors.toList()); + } /** * 查询系统登陆用户数(个)监控信息列表并封装为多折线ECharts图表数据 @@ -156,7 +184,7 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO @Override public Map systemUserNumEcharts(InitialSystemOtherCollectData initialSystemOtherCollectData) { initialSystemOtherCollectData.setCollectType(ServerLogoEnum.登录用户数.getCode()); - List list = initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(initialSystemOtherCollectData); + List list = getOtherMsgSharding(initialSystemOtherCollectData); Map> extractors = new LinkedHashMap<>(); extractors.put("usersNumData", InitialSystemOtherCollectData::getCollectValue); @@ -172,7 +200,7 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO @Override public Map systemSwapSizeFreeEcharts(InitialSystemOtherCollectData initialSystemOtherCollectData) { initialSystemOtherCollectData.setCollectType(ServerLogoEnum.交换卷文件的可用空间.getCode()); - List list = initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(initialSystemOtherCollectData); + List list = getOtherMsgSharding(initialSystemOtherCollectData); Map> extractors = new LinkedHashMap<>(); extractors.put("swapSizeFreeData", InitialSystemOtherCollectData::getCollectValue); @@ -188,7 +216,7 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO @Override public Map memoryUtilizationEcharts(InitialSystemOtherCollectData initialSystemOtherCollectData) { initialSystemOtherCollectData.setCollectType(ServerLogoEnum.内存利用率.getCode()); - List list = initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(initialSystemOtherCollectData); + List list = getOtherMsgSharding(initialSystemOtherCollectData); Map> extractors = new LinkedHashMap<>(); extractors.put("memoryUtilizationData", info -> UnitChangeUtil.formatDecimal(info.getCollectValue())); @@ -204,7 +232,7 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO @Override public Map systemSwapSizePercentEcharts(InitialSystemOtherCollectData initialSystemOtherCollectData) { initialSystemOtherCollectData.setCollectType(ServerLogoEnum.可用交换空间百分比.getCode()); - List list = initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(initialSystemOtherCollectData); + List list = getOtherMsgSharding(initialSystemOtherCollectData); Map> extractors = new LinkedHashMap<>(); extractors.put("swapSizePercentData", info -> UnitChangeUtil.formatDecimal(info.getCollectValue())); @@ -220,7 +248,7 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO @Override public Map memorySizeAvailableEcharts(InitialSystemOtherCollectData initialSystemOtherCollectData) { initialSystemOtherCollectData.setCollectType(ServerLogoEnum.可用内存.getCode()); - List list = initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(initialSystemOtherCollectData); + List list = getOtherMsgSharding(initialSystemOtherCollectData); Map> extractors = new LinkedHashMap<>(); extractors.put("memorySizeAvailableData", data -> { @@ -243,7 +271,7 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO @Override public Map memorySizePercentEcharts(InitialSystemOtherCollectData initialSystemOtherCollectData) { initialSystemOtherCollectData.setCollectType(ServerLogoEnum.可用内存百分比.getCode()); - List list = initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(initialSystemOtherCollectData); + List list = getOtherMsgSharding(initialSystemOtherCollectData); Map> extractors = new LinkedHashMap<>(); extractors.put("memorySizePercentData", info -> UnitChangeUtil.formatDecimal(info.getCollectValue())); @@ -261,11 +289,11 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO public Map procNumEcharts(InitialSystemOtherCollectData initialSystemOtherCollectData) { // 查询总进程数数据 initialSystemOtherCollectData.setCollectType(ServerLogoEnum.进程数.getCode()); - List procNumList = initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(initialSystemOtherCollectData); + List procNumList = getOtherMsgSharding(initialSystemOtherCollectData); // 查询正在运行的进程数数据 initialSystemOtherCollectData.setCollectType(ServerLogoEnum.正在运行的进程数.getCode()); - List procNumRunList = initialSystemOtherCollectDataMapper.selectInitialSystemOtherCollectDataList(initialSystemOtherCollectData); + List procNumRunList = getOtherMsgSharding(initialSystemOtherCollectData); // 先按时间排序 procNumList.sort(Comparator.comparing(InitialSystemOtherCollectData::getCreateTime)); @@ -322,6 +350,8 @@ public class InitialSystemOtherCollectDataServiceImpl implements IInitialSystemO @Override public int deleteInitialSystemOtherCollectData(InitialSystemOtherCollectData deleteData) { + String tableName = TableSubUtil.getTableName(deleteData.getCreateTime(), TABLE_PREFIX); + deleteData.setTableName(tableName); return initialSystemOtherCollectDataMapper.deleteInitialSystemOtherCollectData(deleteData); } } 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 d5f7659..db94af7 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 @@ -268,15 +268,15 @@ public class ProcessSwitchCollectDataService { private void handleSwitchModuleMessage(CollectDataVo switchDataVo, String clientId){ List moduleList = SwitchJsonDataParser.parseJsonData(switchDataVo.getValue(), InitialSwitchOpticalModule.class); if (!moduleList.isEmpty()){ + // 时间戳转换 + long timestamp = switchDataVo.getTimestamp(); + long millis = timestamp * 1000; + Date createTime = new Date(millis / 1000 * 1000); // 去除毫秒 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); + initialSwitchOpticalModuleService.batchInitialSwitchOpticalModule(moduleList, createTime); } } diff --git a/tongran-rocketmq/src/main/resources/mapper/rocketmq/AllMoudleNameMapper.xml b/tongran-rocketmq/src/main/resources/mapper/rocketmq/AllMoudleNameMapper.xml new file mode 100644 index 0000000..b57a07c --- /dev/null +++ b/tongran-rocketmq/src/main/resources/mapper/rocketmq/AllMoudleNameMapper.xml @@ -0,0 +1,82 @@ + + + + + + + + + + + + + + select id, client_id, name, create_time, update_time from all_moudle_name + + + + + + + + insert into all_moudle_name + + client_id, + name, + create_time, + update_time, + + + #{clientId}, + #{name}, + #{createTime}, + #{updateTime}, + + + + + update all_moudle_name + + client_id = #{clientId}, + name = #{name}, + create_time = #{createTime}, + update_time = #{updateTime}, + + where id = #{id} + + + + delete from all_moudle_name where id = #{id} + + + + delete from all_moudle_name where id in + + #{id} + + + + insert IGNORE into all_moudle_name + (client_id, name, create_time, update_time) + values + + ( + #{item.clientId}, + #{item.name}, + #{item.createTime}, + #{item.updateTime} + ) + + + \ No newline at end of file diff --git a/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialSwitchOpticalModuleMapper.xml b/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialSwitchOpticalModuleMapper.xml index a2d8e8d..706a8d3 100644 --- a/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialSwitchOpticalModuleMapper.xml +++ b/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialSwitchOpticalModuleMapper.xml @@ -24,7 +24,7 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" - + select id, fiber_ent_index, fiber_port_name, hw_entity_optical_tx_low_threshold, hw_entity_optical_rx_low_threshold, hw_entity_optical_rx_power, hw_entity_optical_tx_power, client_id, create_time, update_time, create_by, update_by from ${tableName} and fiber_ent_index = #{fiberEntIndex} and fiber_port_name like concat('%', #{fiberPortName}, '%') @@ -147,17 +147,11 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" order by create_time desc limit 1 \ No newline at end of file diff --git a/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialSystemOtherCollectDataMapper.xml b/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialSystemOtherCollectDataMapper.xml index d7ceede..ab36994 100644 --- a/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialSystemOtherCollectDataMapper.xml +++ b/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialSystemOtherCollectDataMapper.xml @@ -20,7 +20,7 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" - insert IGNORE into initial_system_other_collect_data + insert IGNORE into ${tableName} client_id, collect_type, @@ -98,10 +98,10 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" MAX(CASE WHEN t.collect_type = 'memorySizeAvailableCollect' THEN t.collect_value END) as memorySizeAvailableCollect, MAX(CASE WHEN t.collect_type = 'memorySizePercentCollect' THEN t.collect_value END) as memorySizePercentCollect, MAX(CASE WHEN t.collect_type = 'systemUptimeCollect' THEN t.collect_value END) as systemUptimeCollect - FROM `initial_system_other_collect_data` t + FROM ${tableName} t INNER JOIN ( SELECT collect_type, MAX(create_time) as max_time - FROM `initial_system_other_collect_data` + FROM ${tableName} WHERE client_id = #{clientId} GROUP BY collect_type ) latest ON t.collect_type = latest.collect_type AND t.create_time = latest.max_time @@ -109,7 +109,7 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" - delete from initial_system_other_collect_data + delete from ${tableName}