交换机流量数据表改为分表

This commit is contained in:
gaoyutao
2025-11-21 18:32:03 +08:00
parent be2d123902
commit c6e778c74e
19 changed files with 747 additions and 73 deletions
@@ -7,6 +7,7 @@ import com.ruoyi.common.core.web.domain.BaseEntity;
import lombok.Data;
import java.math.BigDecimal;
import java.util.List;
/**
* 交换机流量监控信息对象 initial_switch_info
@@ -69,6 +70,10 @@ public class InitialSwitchInfo extends BaseEntity
private String calculationMode;
/* 单位 */
private String unit;
/** 表名 */
private String tableName;
/** 批量新增集合 */
private List<InitialSwitchInfo> list;
}
@@ -70,4 +70,19 @@ public interface InitialSwitchInfoMapper
public int batchInsertInitialSwitchInfo(@Param("list") List<InitialSwitchInfo> list);
InitialSwitchInfo getSwitchNetDetailsMsg(InitialSwitchInfo initialSwitchInfo);
InitialSwitchInfo getSwitchNetDetailsMsgSharding(InitialSwitchInfo initialSwitchInfo);
/**
* 分表批量插入
* @param initialSwitchInfo
* @return
*/
int batchInsertInitialSwitchInfoSharding(InitialSwitchInfo initialSwitchInfo);
/**
* 分表查询交换机信息
* @param condition
* @return
*/
List<InitialSwitchInfo> selectInitialSwitchInfoListSharding(InitialSwitchInfo condition);
}
@@ -2,6 +2,7 @@ package com.ruoyi.rocketmq.service;
import com.ruoyi.rocketmq.domain.InitialSwitchInfo;
import java.util.Date;
import java.util.List;
import java.util.Map;
@@ -68,6 +69,14 @@ public interface IInitialSwitchInfoService
*/
public int batchInsertInitialSwitchInfo(List<InitialSwitchInfo> list);
/**
* 批量新增交换机流量监控信息
*
* @param list 交换机流量监控信息集合
* @return 结果
*/
public int batchInsertInitialSwitchInfoSharding(List<InitialSwitchInfo> list, Date createTime);
/**
* 交换机网络基础信息
* @param initialSwitchInfo
@@ -1,14 +1,12 @@
package com.ruoyi.rocketmq.service.impl;
import com.ruoyi.common.core.utils.ConvertOtherTypeUtil;
import com.ruoyi.common.core.utils.DateUtils;
import com.ruoyi.common.core.utils.EchartsDataUtils;
import com.ruoyi.common.core.utils.SpeedUtils;
import com.ruoyi.common.core.utils.*;
import com.ruoyi.rocketmq.domain.InitialSwitchInfo;
import com.ruoyi.rocketmq.mapper.InitialSwitchInfoMapper;
import com.ruoyi.rocketmq.service.IInitialSwitchInfoService;
import com.ruoyi.system.api.RemoteRevenueConfigService;
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;
@@ -16,11 +14,9 @@ import org.springframework.transaction.annotation.Transactional;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.util.HashMap;
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业务层处理
@@ -37,6 +33,8 @@ public class InitialSwitchInfoServiceImpl implements IInitialSwitchInfoService
@Autowired
private RemoteRevenueConfigService remoteRevenueConfigService;
private final static String TABLE_PREFIX = "initial_switch_info";
/**
* 查询交换机流量监控信息
*
@@ -127,6 +125,73 @@ public class InitialSwitchInfoServiceImpl implements IInitialSwitchInfoService
throw new RuntimeException("批量保存失败",e);
}
}
/**
* 批量新增交换机流量监控信息
*
* @param list 交换机流量监控信息集合
* @return 结果
*/
@Override
@Transactional(rollbackFor = Exception.class, isolation = Isolation.READ_COMMITTED)
public int batchInsertInitialSwitchInfoSharding(List<InitialSwitchInfo> list, Date createTime) {
if (list == null || list.isEmpty()) {
return 0;
}
// 按表名分组批量插入
Map<String, List<InitialSwitchInfo>> groupedData = list.stream()
.map(data -> {
try {
InitialSwitchInfo processed = new InitialSwitchInfo();
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(
InitialSwitchInfo::getTableName,
LinkedHashMap::new, // 保持插入顺序
Collectors.toList()));
groupedData.forEach((tableName, dataList) -> {
try {
InitialSwitchInfo data = new InitialSwitchInfo();
data.setTableName(tableName);
data.setList(dataList);
initialSwitchInfoMapper.batchInsertInitialSwitchInfoSharding(data);
} catch (Exception e) {
log.error("表{}插入失败", tableName, e);
throw new RuntimeException("批量插入失败", e);
}
});
return 1;
}
/**
* 分表查询交换机网卡流量信息
* @param queryParam
* @return
*/
public List<InitialSwitchInfo> getSwitchMsgSharding(InitialSwitchInfo queryParam) {
// 获取涉及的表名
Set<String> tableNames = TableSubUtil.getTableNamesBetween(queryParam.getStartTime(), queryParam.getEndTime(), TABLE_PREFIX);
// 并行查询各表
return tableNames.parallelStream()
.flatMap(tableName -> {
InitialSwitchInfo condition = new InitialSwitchInfo();
condition.setTableName(tableName);
condition.setClientId(queryParam.getClientId());
condition.setStartTime(queryParam.getStartTime());
condition.setEndTime(queryParam.getEndTime());
return initialSwitchInfoMapper.selectInitialSwitchInfoListSharding(condition).stream();
})
.collect(Collectors.toList());
}
/**
* 获取交换机网络端口的基础信息
@@ -135,7 +200,9 @@ public class InitialSwitchInfoServiceImpl implements IInitialSwitchInfoService
*/
@Override
public InitialSwitchInfo getSwitchNetDetailsMsg(InitialSwitchInfo initialSwitchInfo) {
InitialSwitchInfo info = initialSwitchInfoMapper.getSwitchNetDetailsMsg(initialSwitchInfo);
String tableName = TableSubUtil.getTableName(DateUtils.getNowDate(), TABLE_PREFIX);
initialSwitchInfo.setTableName(tableName);
InitialSwitchInfo info = initialSwitchInfoMapper.getSwitchNetDetailsMsgSharding(initialSwitchInfo);
if(info != null && info.getType()!=null){
info.setType(ConvertOtherTypeUtil.getInterfaceTypeName(Integer.parseInt(info.getType())));
}
@@ -149,7 +216,7 @@ public class InitialSwitchInfoServiceImpl implements IInitialSwitchInfoService
*/
@Override
public Map<String, Object> switchNetSpeedEcharts(InitialSwitchInfo initialSwitchInfo) {
List<InitialSwitchInfo> list = initialSwitchInfoMapper.selectInitialSwitchInfoList(initialSwitchInfo);
List<InitialSwitchInfo> list = getSwitchMsgSharding(initialSwitchInfo);
try {
String unit = SpeedUtils.calculateUnit(list, "inSpeed", "outSpeed");
if(initialSwitchInfo.getUnit() != null){
@@ -185,7 +252,7 @@ public class InitialSwitchInfoServiceImpl implements IInitialSwitchInfoService
*/
@Override
public Map<String, Object> switchNetDiscardsEcharts(InitialSwitchInfo initialSwitchInfo) {
List<InitialSwitchInfo> list = initialSwitchInfoMapper.selectInitialSwitchInfoList(initialSwitchInfo);
List<InitialSwitchInfo> list = getSwitchMsgSharding(initialSwitchInfo);
Map<String, Function<InitialSwitchInfo, ?>> extractors = new LinkedHashMap<>();
extractors.put("netInDiscardsData", info -> info.getIfInDiscards());
extractors.put("netOutDiscardsData", info -> info.getIfOutDiscards());
@@ -202,7 +269,7 @@ public class InitialSwitchInfoServiceImpl implements IInitialSwitchInfoService
*/
@Override
public Map<String, Object> switchNetTotalEcharts(InitialSwitchInfo initialSwitchInfo) {
List<InitialSwitchInfo> list = initialSwitchInfoMapper.selectInitialSwitchInfoList(initialSwitchInfo);
List<InitialSwitchInfo> list = getSwitchMsgSharding(initialSwitchInfo);
Map<String, Function<InitialSwitchInfo, ?>> extractors = new LinkedHashMap<>();
extractors.put("netInTotalData", info -> info.getInBytes().multiply(new BigDecimal(8)));
extractors.put("netOutTotalData", info -> info.getOutBytes().multiply(new BigDecimal(8)));
@@ -219,7 +286,7 @@ public class InitialSwitchInfoServiceImpl implements IInitialSwitchInfoService
*/
@Override
public Map<String, Object> switchNetErrDiscardsEcharts(InitialSwitchInfo initialSwitchInfo) {
List<InitialSwitchInfo> list = initialSwitchInfoMapper.selectInitialSwitchInfoList(initialSwitchInfo);
List<InitialSwitchInfo> list = getSwitchMsgSharding(initialSwitchInfo);
Map<String, Function<InitialSwitchInfo, ?>> extractors = new LinkedHashMap<>();
extractors.put("netInErrDiscardsData", info -> info.getIfInErrors());
extractors.put("netOutErrDiscardsData", info -> info.getIfOutErrors());
@@ -213,7 +213,7 @@ public class ProcessSwitchCollectDataService {
// 临时表 用来计算inSpeed outSeppd
initialSwitchInfoTempService.batchInsertInitialSwitchInfoTemp(switchInfos);
// 初始交换机数据入库
initialSwitchInfoService.batchInsertInitialSwitchInfo(switchInfos);
initialSwitchInfoService.batchInsertInitialSwitchInfoSharding(switchInfos, createTime);
// 业务表入库
InitialSwitchInfoDetailsRemote detailsRemote = new InitialSwitchInfoDetailsRemote();
detailsRemote.setClientId(clientId);