diff --git a/tongran-api/tongran-api-system/src/main/java/com/tongran/system/api/RemoteRevenueConfigService.java b/tongran-api/tongran-api-system/src/main/java/com/tongran/system/api/RemoteRevenueConfigService.java index a09708e..612f6ce 100644 --- a/tongran-api/tongran-api-system/src/main/java/com/tongran/system/api/RemoteRevenueConfigService.java +++ b/tongran-api/tongran-api-system/src/main/java/com/tongran/system/api/RemoteRevenueConfigService.java @@ -133,4 +133,11 @@ public interface RemoteRevenueConfigService @GetMapping("/businessScript/inner/{id}") public R getBusinessScriptMsgByScriptId(@PathVariable("id") Long id, @RequestHeader(SecurityConstants.FROM_SOURCE) String source); + /** + * 保存流量数据 + * @param queryParam 流量数据列表 + * @return 操作结果 + */ + @PostMapping("/revenueConfig/autoSaveServiceRecoverTrafficData") + public R autoSaveServiceRecoverTrafficData(@RequestBody EpsInitialTrafficDataRemote queryParam, @RequestHeader(SecurityConstants.FROM_SOURCE) String source); } diff --git a/tongran-api/tongran-api-system/src/main/java/com/tongran/system/api/factory/RemoteRevenueConfigFallbackFactory.java b/tongran-api/tongran-api-system/src/main/java/com/tongran/system/api/factory/RemoteRevenueConfigFallbackFactory.java index 9739a07..c192497 100644 --- a/tongran-api/tongran-api-system/src/main/java/com/tongran/system/api/factory/RemoteRevenueConfigFallbackFactory.java +++ b/tongran-api/tongran-api-system/src/main/java/com/tongran/system/api/factory/RemoteRevenueConfigFallbackFactory.java @@ -96,6 +96,11 @@ public class RemoteRevenueConfigFallbackFactory implements FallbackFactory getBusinessScriptMsgByScriptId(Long id, String source) { return R.fail("获取错误关键词失败:" + throwable.getMessage()); } + + @Override + public R autoSaveServiceRecoverTrafficData(EpsInitialTrafficDataRemote queryParam, String source) { + return R.fail("保存重试流量数据失败:" + throwable.getMessage()); + } }; } } diff --git a/tongran-common/tongran-common-core/src/main/java/com/tongran/common/core/utils/EchartsDataUtils.java b/tongran-common/tongran-common-core/src/main/java/com/tongran/common/core/utils/EchartsDataUtils.java index 02eca38..c1bd48a 100644 --- a/tongran-common/tongran-common-core/src/main/java/com/tongran/common/core/utils/EchartsDataUtils.java +++ b/tongran-common/tongran-common-core/src/main/java/com/tongran/common/core/utils/EchartsDataUtils.java @@ -118,7 +118,6 @@ public class EchartsDataUtils { // 准备X轴和Y轴数据 List xAxisData = new ArrayList<>(); Map yData = new LinkedHashMap<>(); - // 初始化Y轴数据结构 dataExtractors.keySet().forEach(name -> yData.put(name, new ArrayList<>())); @@ -152,7 +151,7 @@ public class EchartsDataUtils { seriesData.add(getDefaultValue(name, fixedPercentile95Value, xAxisData.size()-1, hasRealData)); } else { // 在数据时间范围外 - seriesData.add(getEmptyDataDefaultValue(name, xAxisData.size()-1)); + seriesData.add(getEmptyDataDefaultValue(name, xAxisData.size()-1, hasRealData)); } } } @@ -338,13 +337,23 @@ public class EchartsDataUtils { /** * 获取空数据默认值(用于数据时间范围外的点) */ - private static Object getEmptyDataDefaultValue(String metricName, int timeIndex) { + private static Object getEmptyDataDefaultValue(String metricName, int timeIndex, boolean hasRealData) { // deployDevice特殊处理:始终补空字符串 if ("deployDevice".equals(metricName)) { return ""; } - - return null; + // 智能补全策略 + if (hasRealData) { + // 数据集中有真实数据:所有缺失点都补null + return null; + } else { + // 数据集中没有真实数据:第一个点补0,其他点补null + if (timeIndex == 0) { + return 0; + } else { + return null; + } + } } /** diff --git a/tongran-modules/tongran-system/src/main/java/com/tongran/system/controller/EpsServerRevenueConfigController.java b/tongran-modules/tongran-system/src/main/java/com/tongran/system/controller/EpsServerRevenueConfigController.java index 8fcf463..e8afcf2 100644 --- a/tongran-modules/tongran-system/src/main/java/com/tongran/system/controller/EpsServerRevenueConfigController.java +++ b/tongran-modules/tongran-system/src/main/java/com/tongran/system/controller/EpsServerRevenueConfigController.java @@ -79,5 +79,14 @@ public class EpsServerRevenueConfigController extends BaseController { return epsServerRevenueConfigService.autoSaveServiceTrafficData(epsServerRevenueConfig); } + /** + * 流量相关数据入库 + */ + @InnerAuth + @PostMapping("/autoSaveServiceRecoverTrafficData") + public R autoSaveServiceRecoverTrafficData(@RequestBody EpsServerRevenueConfig epsServerRevenueConfig) + { + return epsServerRevenueConfigService.autoSaveServiceRecoverTrafficData(epsServerRevenueConfig); + } } 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 bebeca8..46b4c37 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 @@ -74,4 +74,6 @@ public interface EpsInitialTrafficDataMapper { void createSwitchOpMdTable(String tableName); void createDiskInfo(String tableName); + + void batchInsertRecoverDetailTraffic(EpsInitialTrafficData batchData); } 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 debf359..19e14a2 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 @@ -36,6 +36,7 @@ public interface EpsInitialTrafficDataService { * @param dataList 流量数据列表 */ void saveBatch(EpsInitialTrafficData dataList); + void saveBatchRecoverTraffic(EpsInitialTrafficData dataList); /** * 查询流量数据 diff --git a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/IEpsServerRevenueConfigService.java b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/IEpsServerRevenueConfigService.java index 46566e7..fb8f9aa 100644 --- a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/IEpsServerRevenueConfigService.java +++ b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/IEpsServerRevenueConfigService.java @@ -67,6 +67,7 @@ public interface IEpsServerRevenueConfigService * @param epsServerRevenueConfig */ R autoSaveServiceTrafficData(EpsServerRevenueConfig epsServerRevenueConfig); + R autoSaveServiceRecoverTrafficData(EpsServerRevenueConfig epsServerRevenueConfig); /** * 当前在线服务器的流量相关的业务数 * @return 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 c0742c0..353e4ab 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 @@ -245,6 +245,73 @@ public class EpsInitialTrafficDataServiceImpl implements EpsInitialTrafficDataSe } }); } + /** + * 批量保存数据到对应分表 + * @param epsInitialTrafficData 流量数据表 + */ + @Override + @Transactional(rollbackFor = Exception.class, isolation = Isolation.READ_COMMITTED) + public void saveBatchRecoverTraffic(EpsInitialTrafficData epsInitialTrafficData) { + if (epsInitialTrafficData == null || epsInitialTrafficData.getDataList().isEmpty()) { + return; + } + + // 内存去重(基于唯一键) + List distinctList = epsInitialTrafficData.getDataList().stream() + .filter(Objects::nonNull) + .map(data -> { + EpsInitialTrafficData processed = new EpsInitialTrafficData(); + BeanUtils.copyProperties(data, processed); + if (data.getCreateTime() == null) { + processed.setCreateTime(DateUtils.getNowDate()); + } + return processed; + }) + .collect(Collectors.collectingAndThen( + // 使用TreeSet按唯一键去重 + Collectors.toCollection(() -> new TreeSet<>(Comparator.comparing( + d -> String.join("|", + d.getClientId(), + d.getMac(), + d.getName(), + d.getCreateTime().toString() + ) + ))), + ArrayList::new + )); + + // 按表名分组 + Map> groupedData = distinctList.stream() + .map(data -> { + data.setTableName(TableRouterUtil.getTableName( + data.getCreateTime().toInstant() + .atZone(ZoneId.systemDefault()) + .toLocalDateTime() + )); + return data; + }) + .collect(Collectors.groupingBy( + EpsInitialTrafficData::getTableName, + LinkedHashMap::new, + Collectors.toList() + )); + + // 分表插入(带冲突降级) + groupedData.forEach((tableName, list) -> { + try { + EpsInitialTrafficData batchData = new EpsInitialTrafficData(); + BeanUtils.copyProperties(epsInitialTrafficData, batchData); + batchData.setTableName(tableName); + batchData.setDataList(list); + + // 优先尝试批量插入 + epsInitialTrafficDataMapper.batchInsertRecoverDetailTraffic(batchData); + } catch (Exception e) { + log.error("表 {} 插入失败", tableName, e); + throw new RuntimeException("数据入库失败", e); + } + }); + } /** * 查询流量数据 */ diff --git a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/EpsServerRevenueConfigServiceImpl.java b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/EpsServerRevenueConfigServiceImpl.java index c00e559..ba4abde 100644 --- a/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/EpsServerRevenueConfigServiceImpl.java +++ b/tongran-modules/tongran-system/src/main/java/com/tongran/system/service/impl/EpsServerRevenueConfigServiceImpl.java @@ -259,6 +259,102 @@ public class EpsServerRevenueConfigServiceImpl implements IEpsServerRevenueConfi return R.fail("数据保存失败:" + e.getMessage() + ",已成功保存" + successCount + "条"); } } + /** + * 保存流量信息 + * @param epsServerRevenueConfig + */ + @Override + public R autoSaveServiceRecoverTrafficData(EpsServerRevenueConfig epsServerRevenueConfig) { + // 查询初始流量数据 + EpsInitialTrafficData epsInitialTrafficData = new EpsInitialTrafficData(); + epsInitialTrafficData.setStartTime(epsServerRevenueConfig.getStartTime()); + epsInitialTrafficData.setEndTime(epsServerRevenueConfig.getEndTime()); + List dataList = epsInitialTrafficDataService.getAllTraficMsg(epsInitialTrafficData); + + if (dataList == null || dataList.isEmpty()) { + return R.ok("没有需要处理的数据"); + } + + List batchList = new ArrayList<>(); + int batchSize = 1000; // 每批处理数量 + int totalCount = 0; + int successCount = 0; + int batchNumber = 0; + + try { + for (EpsInitialTrafficData initialTrafficData : dataList) { + // 根据clientId查询业务名称 + RmResourceRegistration rmResourceRegistration = new RmResourceRegistration(); + rmResourceRegistration.setClientId(initialTrafficData.getClientId()); + List registerLst = rmResourceRegistrationMapper.selectRmResourceRegistrationList(rmResourceRegistration); + + if(registerLst != null && !registerLst.isEmpty()){ + RmResourceRegistration registerMsg = registerLst.get(0); + // 赋值 + if(registerMsg != null){ + String businessName = registerMsg.getBusinessName(); + if(businessName != null){ + initialTrafficData.setBusinessName(businessName); + // 根据业务名称查询业务代码 + EpsBusiness epsBusiness = epsBusinessMapper.selectEpsBusinessByName(businessName); + if(epsBusiness != null){ + initialTrafficData.setBusinessId(epsBusiness.getId()); + } + } + initialTrafficData.setServiceSn(registerMsg.getHardwareSn()); + initialTrafficData.setRevenueMethod("1"); + } + } + // id自增 + initialTrafficData.setId(null); + batchList.add(initialTrafficData); + + // 达到批次大小时保存 + if (batchList.size() >= batchSize) { + batchNumber++; + totalCount += batchList.size(); + epsInitialTrafficData.setDataList(batchList); + epsInitialTrafficDataService.saveBatchRecoverTraffic(epsInitialTrafficData); + log.info("第{}批流量数据批量入库成功,数据量:{}", batchNumber, batchList.size()); + + // 处理接口名称 + processInterfaceNames(batchList); + + successCount += batchList.size(); + // 清空当前批次,准备下一批 + batchList = new ArrayList<>(); + } + } + + // 处理最后一批不足1000条的数据 + if (!batchList.isEmpty()) { + batchNumber++; + totalCount += batchList.size(); + epsInitialTrafficData.setDataList(batchList); + epsInitialTrafficDataService.saveBatchRecoverTraffic(epsInitialTrafficData); + log.info("第{}批流量数据批量入库成功,数据量:{}", batchNumber, batchList.size()); + + // 处理最后一批的接口名称 + processInterfaceNames(batchList); + + successCount += batchList.size(); + } + + log.info("流量数据批量入库完成,总批次数:{},总数据量:{},成功数量:{}", + batchNumber, totalCount, successCount); + + if (successCount == totalCount) { + return R.ok("数据保存成功,共处理" + successCount + "条数据"); + } else { + return R.fail("数据保存部分成功,应处理" + totalCount + "条,实际成功" + successCount + "条"); + } + + } catch (Exception e) { + log.error("流量数据入库失败,已处理批次:{},成功数量:{},当前批次数量:{},错误原因:{}", + batchNumber, successCount, batchList.size(), e.getMessage(), e); + return R.fail("数据保存失败:" + e.getMessage() + ",已成功保存" + successCount + "条"); + } + } /** * 当前在线服务器的流量相关的业务数 * @return 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 e60c070..ddc0312 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 @@ -377,6 +377,91 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" ) + + INSERT INTO ${tableName} ( + id, + `name`, + `mac`, + `status`, + `type`, + ipV4, + `in_dropped`, + `out_dropped`, + `in_speed`, + `out_speed`, + `total_in_speed`, + `total_out_speed`, + `speed`, + `duplex`, + business_id, + business_name, + service_sn, + node_name, + revenue_method, + package_bandwidth, + create_time, + update_time, + create_by, + update_by, + client_id, + ipv4_in_speed, + ipv4_out_speed, + ipv6_in_speed, + ipv6_out_speed, + total_ipv4_in_speed, + total_ipv4_out_speed, + total_ipv6_in_speed, + total_ipv6_out_speed, + ipV6, + ping_dropped + ) VALUES + + ( + #{data.id,jdbcType=BIGINT}, + #{data.name,jdbcType=VARCHAR}, + #{data.mac,jdbcType=VARCHAR}, + #{data.status,jdbcType=VARCHAR}, + #{data.type,jdbcType=VARCHAR}, + #{data.ipV4,jdbcType=VARCHAR}, + #{data.inDropped,jdbcType=DECIMAL}, + #{data.outDropped,jdbcType=DECIMAL}, + #{data.inSpeed,jdbcType=VARCHAR}, + #{data.outSpeed,jdbcType=VARCHAR}, + #{data.totalInSpeed,jdbcType=VARCHAR}, + #{data.totalOutSpeed,jdbcType=VARCHAR}, + #{data.speed,jdbcType=VARCHAR}, + #{data.duplex,jdbcType=VARCHAR}, + #{data.businessId,jdbcType=VARCHAR}, + #{data.businessName,jdbcType=VARCHAR}, + #{data.serviceSn,jdbcType=VARCHAR}, + #{data.nodeName,jdbcType=VARCHAR}, + #{data.revenueMethod,jdbcType=VARCHAR}, + #{data.packageBandwidth,jdbcType=DECIMAL}, + #{data.createTime,jdbcType=TIMESTAMP}, + #{data.updateTime,jdbcType=TIMESTAMP}, + #{data.createBy,jdbcType=VARCHAR}, + #{data.updateBy,jdbcType=VARCHAR}, + #{data.clientId,jdbcType=VARCHAR}, + #{data.ipv4InSpeed,jdbcType=VARCHAR}, + #{data.ipv4OutSpeed,jdbcType=VARCHAR}, + #{data.ipv6InSpeed,jdbcType=VARCHAR}, + #{data.ipv6OutSpeed,jdbcType=VARCHAR}, + #{data.totalIpv4InSpeed,jdbcType=VARCHAR}, + #{data.totalIpv4OutSpeed,jdbcType=VARCHAR}, + #{data.totalIpv6InSpeed,jdbcType=VARCHAR}, + #{data.totalIpv6OutSpeed,jdbcType=VARCHAR}, + #{data.ipV6,jdbcType=VARCHAR}, + #{data.pingDropped,jdbcType=DECIMAL} + ) + + ON DUPLICATE KEY UPDATE + `in_speed` = IF(VALUES(`in_speed`) IS NOT NULL, VALUES(`in_speed`), `in_speed`), + `out_speed` = IF(VALUES(`out_speed`) IS NOT NULL, VALUES(`out_speed`), `out_speed`), + `ipv4_in_speed` = IF(VALUES(`ipv4_in_speed`) IS NOT NULL, VALUES(`ipv4_in_speed`), `ipv4_in_speed`), + `ipv4_out_speed` = IF(VALUES(`ipv4_out_speed`) IS NOT NULL, VALUES(`ipv4_out_speed`), `ipv4_out_speed`), + `ipv6_in_speed` = IF(VALUES(`ipv6_in_speed`) IS NOT NULL, VALUES(`ipv6_in_speed`), `ipv6_in_speed`), + `ipv6_out_speed` = IF(VALUES(`ipv6_out_speed`) IS NOT NULL, VALUES(`ipv6_out_speed`), `ipv6_out_speed`) + select id, @@ -229,6 +255,7 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" and client_id = #{clientId} and create_time = #{startTime} + and in_speed is not null \ No newline at end of file