From c921a83d27671a2f1dd6c76a302aa99aa9214e68 Mon Sep 17 00:00:00 2001 From: gaoyutao Date: Thu, 15 Jan 2026 10:27:40 +0800 Subject: [PATCH] =?UTF-8?q?=E7=A3=81=E7=9B=98=E8=A1=A8=E5=88=86=E8=A1=A8?= =?UTF-8?q?=EF=BC=8C=E5=A2=9E=E5=8A=A0=E7=A3=81=E7=9B=98=E5=90=8D=E7=A7=B0?= =?UTF-8?q?=E5=AD=98=E5=82=A8CRUD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../controller/CalculateController.java | 6 + .../mapper/EpsInitialTrafficDataMapper.java | 2 + .../service/EpsInitialTrafficDataService.java | 6 + .../EpsInitialTrafficDataServiceImpl.java | 28 +++++ .../system/EpsInitialTrafficDataMapper.xml | 23 ++++ .../controller/AllDiskNameController.java | 98 +++++++++++++++ .../tongran/rocketmq/domain/AllDiskName.java | 116 +++++++++++++++++ .../rocketmq/domain/InitialDiskInfo.java | 15 +++ .../handler/DeviceMessageHandler.java | 2 +- .../rocketmq/handler/MessageHandler.java | 2 +- .../rocketmq/mapper/AllDiskNameMapper.java | 64 ++++++++++ .../mapper/InitialDiskInfoMapper.java | 7 +- .../rocketmq/service/IAllDiskNameService.java | 65 ++++++++++ .../service/IInitialDiskInfoService.java | 4 +- .../service/impl/AllDiskNameServiceImpl.java | 118 +++++++++++++++++ .../impl/InitialDiskInfoServiceImpl.java | 84 +++++++++++-- .../mapper/rocketmq/AllDiskNameMapper.xml | 119 ++++++++++++++++++ .../mapper/rocketmq/InitialDiskInfoMapper.xml | 60 +++++++-- 18 files changed, 788 insertions(+), 31 deletions(-) create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/controller/AllDiskNameController.java create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/AllDiskName.java create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/AllDiskNameMapper.java create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IAllDiskNameService.java create mode 100644 tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/AllDiskNameServiceImpl.java create mode 100644 tongran-rocketmq/src/main/resources/mapper/rocketmq/AllDiskNameMapper.xml diff --git a/tongran-modules/tongran-system/src/main/java/com/tongran/system/controller/CalculateController.java b/tongran-modules/tongran-system/src/main/java/com/tongran/system/controller/CalculateController.java index 3a8d784..73a173e 100644 --- a/tongran-modules/tongran-system/src/main/java/com/tongran/system/controller/CalculateController.java +++ b/tongran-modules/tongran-system/src/main/java/com/tongran/system/controller/CalculateController.java @@ -1,6 +1,7 @@ package com.tongran.system.controller; import com.tongran.common.core.web.controller.BaseController; +import com.tongran.common.core.web.domain.AjaxResult; import com.tongran.system.domain.EpsInitialTrafficData; import com.tongran.system.domain.InitialSwitchInfoDetails; import com.tongran.system.service.EpsInitialTrafficDataService; @@ -24,6 +25,11 @@ public class CalculateController extends BaseController { @Autowired private IInitialSwitchInfoDetailsService initialSwitchInfoDetailsService; + @GetMapping("/createTables") + public AjaxResult createTables(Long plusMonth){ + epsInitialTrafficDataService.createNextMonthTables(plusMonth); + return success(); + } @GetMapping("/calculate95BandwidthDaily") public void calculate95BandwidthDaily(String day){ // 获取昨天的日期范围(北京时间) 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 70c14bb..bebeca8 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 @@ -72,4 +72,6 @@ public interface EpsInitialTrafficDataMapper { void createOtherMsgTable(String tableName); void createSwitchOpMdTable(String tableName); + + void createDiskInfo(String tableName); } diff --git a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/EpsInitialTrafficDataService.java b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/EpsInitialTrafficDataService.java index 38ccdf7..debf359 100644 --- a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/EpsInitialTrafficDataService.java +++ b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/EpsInitialTrafficDataService.java @@ -19,6 +19,12 @@ public interface EpsInitialTrafficDataService { * 每月25日0点自动执行 */ void createNextMonthTables(); + + /** + * 创建分表 + * @param plusMonth 下几个月 + */ + void createNextMonthTables(Long plusMonth); /** * 保存单条流量数据 * @param data 流量数据 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 6c6652d..4b6524d 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 @@ -74,8 +74,31 @@ public class EpsInitialTrafficDataServiceImpl implements EpsInitialTrafficDataSe createOtherMsgTable(year, month); // 创建交换机光模块表 createSwitchOpMdTable(year, month); + // 创建磁盘信息表 + createDiskInfo(year, month); } + public void createNextMonthTables(Long plusMonth) { + LocalDate nextMonth = LocalDate.now().plusMonths(plusMonth); + int year = nextMonth.getYear(); + int month = nextMonth.getMonthValue(); + // 创建流量详情表 + createTrafficDetailsTable(year, month); + // 创建流量初始表 + createTrafficStatsTable(year, month); + // 创建mtr探测丢包结果表 + createMtrProbeResultTable(year, month); + // 创建交换机流量初始表 + createSwitchInfoTable(year, month); + // 创建交换机流量业务表 + createSwitchInfoDetailsTable(year, month); + // 创建服务器其他信息表 + createOtherMsgTable(year, month); + // 创建交换机光模块表 + createSwitchOpMdTable(year, month); + // 创建磁盘信息表 + createDiskInfo(year, month); + } private void createTrafficDetailsTable(int year, int month) { createRangeTables(year, month, "eps_traffic_details", (tableName) -> { epsInitialTrafficDataMapper.createEpsTrafficTable(tableName); @@ -114,6 +137,11 @@ public class EpsInitialTrafficDataServiceImpl implements EpsInitialTrafficDataSe epsInitialTrafficDataMapper.createSwitchOpMdTable(tableName); }); } + private void createDiskInfo(int year, int month) { + createRangeTables(year, month, "initial_disk_info", (tableName) -> { + epsInitialTrafficDataMapper.createDiskInfo(tableName); + }); + } /** 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 12ea429..dfe478a 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 @@ -201,6 +201,29 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" UNIQUE KEY uk_client_fiber_time (client_id, create_time, fiber_port_name) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci COMMENT='光模块信息表'; + + CREATE TABLE IF NOT EXISTS ${tableName} ( + `id` bigint(20) NOT NULL AUTO_INCREMENT COMMENT '主键ID', + `client_id` varchar(255) NOT NULL COMMENT '客户端ID', + `name` varchar(100) NOT NULL COMMENT '磁盘名称(如sda、sdb等)', + `serial` varchar(100) COMMENT '磁盘序列号', + `total` bigint(20) COMMENT '磁盘总大小(GB)', + `write_speed` bigint(20) COMMENT '磁盘写入速率(字节/秒)', + `read_speed` bigint(20) COMMENT '磁盘读取速率(字节/秒)', + `write_times` bigint(20) COMMENT '磁盘写入次数', + `read_times` bigint(20) COMMENT '磁盘读取次数', + `write_bytes` bigint(20) COMMENT '磁盘写入总字节数', + `read_bytes` bigint(20) COMMENT '磁盘读取总字节数', + `create_by` varchar(50) COMMENT '创建人', + `update_by` varchar(50) COMMENT '更新人', + `create_time` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', + `update_time` datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间', + `type` varchar(255) COMMENT '磁盘类型', + `used_space` bigint(20) COMMENT '已用空间', + PRIMARY KEY (`id`), + UNIQUE KEY uk_client_disk_time (`client_id`, `name`, `create_time`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci COMMENT='磁盘监控信息表'; + INSERT INTO ${tableName} ( diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/controller/AllDiskNameController.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/controller/AllDiskNameController.java new file mode 100644 index 0000000..754c460 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/controller/AllDiskNameController.java @@ -0,0 +1,98 @@ +package com.tongran.rocketmq.controller; + +import com.tongran.common.core.utils.poi.ExcelUtil; +import com.tongran.common.core.web.controller.BaseController; +import com.tongran.common.core.web.domain.AjaxResult; +import com.tongran.common.core.web.page.TableDataInfo; +import com.tongran.common.log.annotation.Log; +import com.tongran.common.log.enums.BusinessType; +import com.tongran.common.security.annotation.RequiresPermissions; +import com.tongran.rocketmq.domain.AllDiskName; +import com.tongran.rocketmq.service.IAllDiskNameService; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.web.bind.annotation.*; + +import javax.servlet.http.HttpServletResponse; +import java.util.List; + +/** + * 磁盘名称存储Controller + * + * @author tongran + * @date 2026-01-14 + */ +@RestController +@RequestMapping("/allDiskName") +public class AllDiskNameController extends BaseController +{ + @Autowired + private IAllDiskNameService allDiskNameService; + + /** + * 查询磁盘名称存储列表 + */ + @RequiresPermissions("rocketmq:allDiskName:list") + @GetMapping("/list") + public TableDataInfo list(AllDiskName allDiskName) + { + startPage(); + List list = allDiskNameService.selectAllDiskNameList(allDiskName); + return getDataTable(list); + } + + /** + * 导出磁盘名称存储列表 + */ + @RequiresPermissions("rocketmq:allDiskName:export") + @Log(title = "磁盘名称存储", businessType = BusinessType.EXPORT) + @PostMapping("/export") + public void export(HttpServletResponse response, AllDiskName allDiskName) + { + List list = allDiskNameService.selectAllDiskNameList(allDiskName); + ExcelUtil util = new ExcelUtil(AllDiskName.class); + util.exportExcel(response, list, "磁盘名称存储数据"); + } + + /** + * 获取磁盘名称存储详细信息 + */ + @RequiresPermissions("rocketmq:allDiskName:query") + @GetMapping(value = "/{id}") + public AjaxResult getInfo(@PathVariable("id") Long id) + { + return success(allDiskNameService.selectAllDiskNameById(id)); + } + + /** + * 新增磁盘名称存储 + */ + @RequiresPermissions("rocketmq:allDiskName:add") + @Log(title = "磁盘名称存储", businessType = BusinessType.INSERT) + @PostMapping + public AjaxResult add(@RequestBody AllDiskName allDiskName) + { + return toAjax(allDiskNameService.insertAllDiskName(allDiskName)); + } + + /** + * 修改磁盘名称存储 + */ + @RequiresPermissions("rocketmq:allDiskName:edit") + @Log(title = "磁盘名称存储", businessType = BusinessType.UPDATE) + @PutMapping + public AjaxResult edit(@RequestBody AllDiskName allDiskName) + { + return toAjax(allDiskNameService.updateAllDiskName(allDiskName)); + } + + /** + * 删除磁盘名称存储 + */ + @RequiresPermissions("rocketmq:allDiskName:remove") + @Log(title = "磁盘名称存储", businessType = BusinessType.DELETE) + @DeleteMapping("/{ids}") + public AjaxResult remove(@PathVariable Long[] ids) + { + return toAjax(allDiskNameService.deleteAllDiskNameByIds(ids)); + } +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/AllDiskName.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/AllDiskName.java new file mode 100644 index 0000000..03afaf9 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/AllDiskName.java @@ -0,0 +1,116 @@ +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_disk_name + * + * @author tongran + * @date 2026-01-14 + */ +public class AllDiskName extends BaseEntity +{ + private static final long serialVersionUID = 1L; + + /** 主键ID */ + private Long id; + + /** 客户端ID */ + @Excel(name = "客户端ID") + private String clientId; + + /** 磁盘名称 */ + @Excel(name = "磁盘名称") + private String name; + + /** 磁盘状态(0:丢失,1:存在) */ + @Excel(name = "磁盘状态(0:丢失,1:存在)") + private Integer status; + + /** 读取IOPS */ + @Excel(name = "读取IOPS") + private String readIops; + + /** 写入IOPS */ + @Excel(name = "写入IOPS") + private String writeIops; + + 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; + } + + public void setStatus(Integer status) + { + this.status = status; + } + + public Integer getStatus() + { + return status; + } + + public void setReadIops(String readIops) + { + this.readIops = readIops; + } + + public String getReadIops() + { + return readIops; + } + + public void setWriteIops(String writeIops) + { + this.writeIops = writeIops; + } + + public String getWriteIops() + { + return writeIops; + } + + @Override + public String toString() { + return new ToStringBuilder(this,ToStringStyle.MULTI_LINE_STYLE) + .append("id", getId()) + .append("clientId", getClientId()) + .append("name", getName()) + .append("status", getStatus()) + .append("readIops", getReadIops()) + .append("writeIops", getWriteIops()) + .append("createTime", getCreateTime()) + .append("updateTime", getUpdateTime()) + .append("createBy", getCreateBy()) + .append("updateBy", getUpdateBy()) + .toString(); + } +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialDiskInfo.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialDiskInfo.java index dbea8db..ce9038d 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialDiskInfo.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/domain/InitialDiskInfo.java @@ -4,6 +4,8 @@ import com.tongran.common.core.annotation.Excel; import com.tongran.common.core.web.domain.BaseEntity; import lombok.Data; +import java.util.List; + /** * 磁盘监控信息对象 initial_disk_info * @@ -69,5 +71,18 @@ public class InitialDiskInfo extends BaseEntity private String readBytesStr; /** 换算后的写入次数 */ private String writeBytesStr; + /** 磁盘类型 */ + private String type; + /** 已用空间 */ + private Long usedSpace; + /** 读IOPS */ + private Long readIops; + /** 写IOPS */ + private Long writeIops; + /** 表名 */ + private String tableName; + /** 批量插入列表 */ + private List list; + } 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 a97a290..21f283f 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 @@ -265,7 +265,7 @@ public class DeviceMessageHandler { iface.setCreateTime(createTime); }); // 初始磁盘数据入库 - initialDiskInfoService.batchInsertInitialDiskInfo(disks); + initialDiskInfoService.batchInsertInitialDiskInfo(disks, createTime); }else{ throw new RuntimeException("磁盘data数据为空"); } 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 5aa27f2..a6d8461 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 @@ -550,7 +550,7 @@ public class MessageHandler { iface.setCreateTime(createTime); }); // 初始磁盘数据入库 - initialDiskInfoService.batchInsertInitialDiskInfo(disks); + initialDiskInfoService.batchInsertInitialDiskInfo(disks, createTime); }else{ throw new RuntimeException("磁盘data数据为空"); } diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/AllDiskNameMapper.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/AllDiskNameMapper.java new file mode 100644 index 0000000..1083037 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/AllDiskNameMapper.java @@ -0,0 +1,64 @@ +package com.tongran.rocketmq.mapper; + +import com.tongran.rocketmq.domain.AllDiskName; + +import java.util.List; + +/** + * 磁盘名称存储Mapper接口 + * + * @author tongran + * @date 2026-01-14 + */ +public interface AllDiskNameMapper +{ + /** + * 查询磁盘名称存储 + * + * @param id 磁盘名称存储主键 + * @return 磁盘名称存储 + */ + public AllDiskName selectAllDiskNameById(Long id); + + /** + * 查询磁盘名称存储列表 + * + * @param allDiskName 磁盘名称存储 + * @return 磁盘名称存储集合 + */ + public List selectAllDiskNameList(AllDiskName allDiskName); + + /** + * 新增磁盘名称存储 + * + * @param allDiskName 磁盘名称存储 + * @return 结果 + */ + public int insertAllDiskName(AllDiskName allDiskName); + + /** + * 修改磁盘名称存储 + * + * @param allDiskName 磁盘名称存储 + * @return 结果 + */ + public int updateAllDiskName(AllDiskName allDiskName); + + /** + * 删除磁盘名称存储 + * + * @param id 磁盘名称存储主键 + * @return 结果 + */ + public int deleteAllDiskNameById(Long id); + + /** + * 批量删除磁盘名称存储 + * + * @param ids 需要删除的数据主键集合 + * @return 结果 + */ + public int deleteAllDiskNameByIds(Long[] ids); + + int batchInsertAllDistName(List dataList); +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/InitialDiskInfoMapper.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/InitialDiskInfoMapper.java index 0dd5c94..d4aa682 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/InitialDiskInfoMapper.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/mapper/InitialDiskInfoMapper.java @@ -1,7 +1,6 @@ package com.tongran.rocketmq.mapper; import com.tongran.rocketmq.domain.InitialDiskInfo; -import org.springframework.data.repository.query.Param; import java.util.List; import java.util.Map; @@ -65,10 +64,10 @@ public interface InitialDiskInfoMapper /** * 批量新增磁盘监控信息 * - * @param list 磁盘监控信息集合 + * @param initialDiskInfo 磁盘监控信息 * @return 结果 */ - public int batchInsertInitialDiskInfo(@Param("list") List list); + public int batchInsertInitialDiskInfo(InitialDiskInfo initialDiskInfo); /** * 获取磁盘设备基础信息 @@ -83,4 +82,6 @@ public interface InitialDiskInfoMapper * @return */ List getAllDistName(InitialDiskInfo initialDiskInfo); + + List selectInitialDiskInfoListByCondition(InitialDiskInfo condition); } diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IAllDiskNameService.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IAllDiskNameService.java new file mode 100644 index 0000000..24c9fe5 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IAllDiskNameService.java @@ -0,0 +1,65 @@ +package com.tongran.rocketmq.service; + +import com.tongran.rocketmq.domain.AllDiskName; +import com.tongran.rocketmq.domain.InitialDiskInfo; + +import java.util.List; + +/** + * 磁盘名称存储Service接口 + * + * @author tongran + * @date 2026-01-14 + */ +public interface IAllDiskNameService +{ + /** + * 查询磁盘名称存储 + * + * @param id 磁盘名称存储主键 + * @return 磁盘名称存储 + */ + public AllDiskName selectAllDiskNameById(Long id); + + /** + * 查询磁盘名称存储列表 + * + * @param allDiskName 磁盘名称存储 + * @return 磁盘名称存储集合 + */ + public List selectAllDiskNameList(AllDiskName allDiskName); + + /** + * 新增磁盘名称存储 + * + * @param allDiskName 磁盘名称存储 + * @return 结果 + */ + public int insertAllDiskName(AllDiskName allDiskName); + + /** + * 修改磁盘名称存储 + * + * @param allDiskName 磁盘名称存储 + * @return 结果 + */ + public int updateAllDiskName(AllDiskName allDiskName); + + /** + * 批量删除磁盘名称存储 + * + * @param ids 需要删除的磁盘名称存储主键集合 + * @return 结果 + */ + public int deleteAllDiskNameByIds(Long[] ids); + + /** + * 删除磁盘名称存储信息 + * + * @param id 磁盘名称存储主键 + * @return 结果 + */ + public int deleteAllDiskNameById(Long id); + + int batchInsertAllDistName(List dataList); +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IInitialDiskInfoService.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IInitialDiskInfoService.java index 51443ff..9e7a6f6 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IInitialDiskInfoService.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/IInitialDiskInfoService.java @@ -2,6 +2,7 @@ package com.tongran.rocketmq.service; import com.tongran.rocketmq.domain.InitialDiskInfo; +import java.util.Date; import java.util.List; import java.util.Map; @@ -64,9 +65,10 @@ public interface IInitialDiskInfoService * 批量新增磁盘监控信息 * * @param list 磁盘监控信息集合 + * @param createTime 采集时间 * @return 结果 */ - public int batchInsertInitialDiskInfo(List list); + public int batchInsertInitialDiskInfo(List list, Date createTime); /** * 磁盘设备/dev/sda基础信息 * @param initialDiskInfo diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/AllDiskNameServiceImpl.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/AllDiskNameServiceImpl.java new file mode 100644 index 0000000..e100724 --- /dev/null +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/AllDiskNameServiceImpl.java @@ -0,0 +1,118 @@ +package com.tongran.rocketmq.service.impl; + +import com.tongran.common.core.utils.DateUtils; +import com.tongran.rocketmq.domain.AllDiskName; +import com.tongran.rocketmq.domain.InitialDiskInfo; +import com.tongran.rocketmq.mapper.AllDiskNameMapper; +import com.tongran.rocketmq.service.IAllDiskNameService; +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-14 + */ +@Service +public class AllDiskNameServiceImpl implements IAllDiskNameService +{ + @Autowired + private AllDiskNameMapper allDiskNameMapper; + + /** + * 查询磁盘名称存储 + * + * @param id 磁盘名称存储主键 + * @return 磁盘名称存储 + */ + @Override + public AllDiskName selectAllDiskNameById(Long id) + { + return allDiskNameMapper.selectAllDiskNameById(id); + } + + /** + * 查询磁盘名称存储列表 + * + * @param allDiskName 磁盘名称存储 + * @return 磁盘名称存储 + */ + @Override + public List selectAllDiskNameList(AllDiskName allDiskName) + { + return allDiskNameMapper.selectAllDiskNameList(allDiskName); + } + + /** + * 新增磁盘名称存储 + * + * @param allDiskName 磁盘名称存储 + * @return 结果 + */ + @Override + public int insertAllDiskName(AllDiskName allDiskName) + { + allDiskName.setCreateTime(DateUtils.getNowDate()); + return allDiskNameMapper.insertAllDiskName(allDiskName); + } + + /** + * 修改磁盘名称存储 + * + * @param allDiskName 磁盘名称存储 + * @return 结果 + */ + @Override + public int updateAllDiskName(AllDiskName allDiskName) + { + allDiskName.setUpdateTime(DateUtils.getNowDate()); + return allDiskNameMapper.updateAllDiskName(allDiskName); + } + + /** + * 批量删除磁盘名称存储 + * + * @param ids 需要删除的磁盘名称存储主键 + * @return 结果 + */ + @Override + public int deleteAllDiskNameByIds(Long[] ids) + { + return allDiskNameMapper.deleteAllDiskNameByIds(ids); + } + + /** + * 删除磁盘名称存储信息 + * + * @param id 磁盘名称存储主键 + * @return 结果 + */ + @Override + public int deleteAllDiskNameById(Long id) + { + return allDiskNameMapper.deleteAllDiskNameById(id); + } + + @Override + public int batchInsertAllDistName(List dataList) { + if(dataList == null){ + dataList = new ArrayList<>(); + } + List nameList = new ArrayList<>(); + for (InitialDiskInfo initialDiskInfo : dataList) { + AllDiskName allDiskName = new AllDiskName(); + allDiskName.setClientId(initialDiskInfo.getClientId()); + allDiskName.setName(initialDiskInfo.getName()); + allDiskName.setStatus(1); + allDiskName.setCreateTime(initialDiskInfo.getCreateTime()); + allDiskName.setUpdateTime(DateUtils.getNowDate()); + nameList.add(allDiskName); + } + allDiskNameMapper.batchInsertAllDistName(nameList); + return 1; + } +} diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialDiskInfoServiceImpl.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialDiskInfoServiceImpl.java index 9d8c434..2b8d6ae 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialDiskInfoServiceImpl.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/service/impl/InitialDiskInfoServiceImpl.java @@ -2,20 +2,22 @@ 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.InitialDiskInfo; import com.tongran.rocketmq.mapper.InitialDiskInfoMapper; +import com.tongran.rocketmq.service.IAllDiskNameService; import com.tongran.rocketmq.service.IInitialDiskInfoService; 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业务层处理 @@ -29,6 +31,9 @@ public class InitialDiskInfoServiceImpl implements IInitialDiskInfoService { @Autowired private InitialDiskInfoMapper initialDiskInfoMapper; + @Autowired + private IAllDiskNameService allDiskNameService; + private final static String TABLE_PREFIX = "initial_disk_info"; /** * 查询磁盘监控信息 @@ -111,13 +116,44 @@ public class InitialDiskInfoServiceImpl implements IInitialDiskInfoService */ @Override @Transactional(rollbackFor = Exception.class, isolation = Isolation.READ_COMMITTED) - public int batchInsertInitialDiskInfo(List list) { - try { - return initialDiskInfoMapper.batchInsertInitialDiskInfo(list); - }catch (Exception e){ - log.error("批量插入磁盘信息失败,失败数量:{}", list.size(), e); - throw new RuntimeException("批量保存失败",e); + public int batchInsertInitialDiskInfo(List list, Date createTime) { + if (list == null || list.isEmpty()) { + return 0; } + // 按表名分组批量插入 + Map> groupedData = list.stream() + .map(data -> { + try { + InitialDiskInfo processed = new InitialDiskInfo(); + 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( + InitialDiskInfo::getTableName, + LinkedHashMap::new, // 保持插入顺序 + Collectors.toList())); + + groupedData.forEach((tableName, dataList) -> { + try { + InitialDiskInfo data = new InitialDiskInfo(); + data.setTableName(tableName); + data.setList(dataList); + initialDiskInfoMapper.batchInsertInitialDiskInfo(data); + // 记录磁盘名称 + allDiskNameService.batchInsertAllDistName(dataList); + } catch (Exception e) { + log.error("表{}插入失败", tableName, e); + throw new RuntimeException("批量插入失败", e); + } + }); + return 1; } /** * 磁盘设备/dev/sda基础信息 @@ -126,6 +162,8 @@ public class InitialDiskInfoServiceImpl implements IInitialDiskInfoService */ @Override public InitialDiskInfo getDistDetailsMsg(InitialDiskInfo initialDiskInfo) { + String tableName = TableSubUtil.getTableName(DateUtils.getNowDate(), TABLE_PREFIX); + initialDiskInfo.setTableName(tableName); InitialDiskInfo info = initialDiskInfoMapper.getDistDetailsMsgByClientId(initialDiskInfo); if(info != null){ long gbUnit = 1024L * 1024 * 1024; @@ -137,6 +175,28 @@ public class InitialDiskInfoServiceImpl implements IInitialDiskInfoService } return info; } + /** + * 分表查询硬盘信息 + * @param queryParam + * @return + */ + public List getDistInfoSharding(InitialDiskInfo queryParam) { + // 获取涉及的表名 + Set tableNames = TableSubUtil.getExistingTableNamesBetween(queryParam.getStartTime(), queryParam.getEndTime(), TABLE_PREFIX); + + // 并行查询各表 + return tableNames.parallelStream() + .flatMap(tableName -> { + InitialDiskInfo condition = new InitialDiskInfo(); + condition.setTableName(tableName); + condition.setClientId(queryParam.getClientId()); + condition.setName(queryParam.getName()); + condition.setStartTime(queryParam.getStartTime()); + condition.setEndTime(queryParam.getEndTime()); + return initialDiskInfoMapper.selectInitialDiskInfoListByCondition(condition).stream(); + }) + .collect(Collectors.toList()); + } /** * /dev/sda读写速率(KB/s) * @param initialDiskInfo @@ -144,7 +204,7 @@ public class InitialDiskInfoServiceImpl implements IInitialDiskInfoService */ @Override public Map rwSpeedEcharts(InitialDiskInfo initialDiskInfo) { - List list = initialDiskInfoMapper.selectInitialDiskInfoList(initialDiskInfo); + List list = getDistInfoSharding(initialDiskInfo); Map> extractors = new LinkedHashMap<>(); extractors.put("readSpeedData", info -> info.getReadSpeed() / 1024.0); extractors.put("writeSpeedData", info -> info.getWriteSpeed() / 1024.0); @@ -157,7 +217,7 @@ public class InitialDiskInfoServiceImpl implements IInitialDiskInfoService */ @Override public Map rwTimesEcharts(InitialDiskInfo initialDiskInfo) { - List list = initialDiskInfoMapper.selectInitialDiskInfoList(initialDiskInfo); + List list = getDistInfoSharding(initialDiskInfo); Map> extractors = new LinkedHashMap<>(); extractors.put("readTimesData", info -> info.getReadTimes()); extractors.put("writeTimesData", info -> info.getWriteTimes()); @@ -170,7 +230,7 @@ public class InitialDiskInfoServiceImpl implements IInitialDiskInfoService */ @Override public Map rwBytesEcharts(InitialDiskInfo initialDiskInfo) { - List list = initialDiskInfoMapper.selectInitialDiskInfoList(initialDiskInfo); + List list = getDistInfoSharding(initialDiskInfo); Map> extractors = new LinkedHashMap<>(); extractors.put("readBytesData", info -> info.getReadBytes()); extractors.put("writeBytesData", info -> info.getWriteBytes()); diff --git a/tongran-rocketmq/src/main/resources/mapper/rocketmq/AllDiskNameMapper.xml b/tongran-rocketmq/src/main/resources/mapper/rocketmq/AllDiskNameMapper.xml new file mode 100644 index 0000000..709d35f --- /dev/null +++ b/tongran-rocketmq/src/main/resources/mapper/rocketmq/AllDiskNameMapper.xml @@ -0,0 +1,119 @@ + + + + + + + + + + + + + + + + + + + select id, client_id, name, status, read_iops, write_iops, create_time, update_time, create_by, update_by from all_disk_name + + + + + + + + insert into all_disk_name + + client_id, + name, + status, + read_iops, + write_iops, + create_time, + update_time, + create_by, + update_by, + + + #{clientId}, + #{name}, + #{status}, + #{readIops}, + #{writeIops}, + #{createTime}, + #{updateTime}, + #{createBy}, + #{updateBy}, + + on duplicate key update + + name = #{name}, + status = #{status}, + read_iops = #{readIops}, + write_iops = #{writeIops}, + update_time = #{updateTime}, + update_by = #{updateBy}, + + + + + update all_disk_name + + client_id = #{clientId}, + name = #{name}, + status = #{status}, + read_iops = #{readIops}, + write_iops = #{writeIops}, + create_time = #{createTime}, + update_time = #{updateTime}, + create_by = #{createBy}, + update_by = #{updateBy}, + + where id = #{id} + + + + delete from all_disk_name where id = #{id} + + + + delete from all_disk_name where id in + + #{id} + + + + insert IGNORE into all_disk_name + (client_id, name, status, read_iops, write_iops, create_time, update_time, create_by, update_by) + values + + ( + #{item.clientId}, + #{item.name}, + #{item.status}, + #{item.readIops}, + #{item.writeIops}, + #{item.createTime}, + #{item.updateTime}, + #{item.createBy}, + #{item.updateBy} + ) + + + \ No newline at end of file diff --git a/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialDiskInfoMapper.xml b/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialDiskInfoMapper.xml index ce0fe37..8ddead9 100644 --- a/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialDiskInfoMapper.xml +++ b/tongran-rocketmq/src/main/resources/mapper/rocketmq/InitialDiskInfoMapper.xml @@ -1,9 +1,9 @@ + PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" + "http://mybatis.org/dtd/mybatis-3-mapper.dtd"> - + @@ -16,6 +16,9 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" + + + @@ -23,20 +26,34 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" - select id, client_id, name, serial, total, write_speed, read_speed, write_times, read_times, write_bytes, read_bytes, create_by, update_by, create_time, update_time from initial_disk_info + select id, client_id, name, serial, total, write_speed, read_speed, write_times, read_times, write_bytes, read_bytes, type, used_space, create_by, update_by, create_time, update_time from initial_disk_info - + + - + select id, client_id, name, serial, total, write_speed, read_speed, write_times, read_times, write_bytes, + read_bytes, type, used_space, create_by, update_by, create_time, update_time + from ${tableName} where client_id = #{clientId} and `name`= #{name} order by create_time desc limit 1 +