优化图形补0、null机制,优化监控看板接口
mtr监控策略相关接口
This commit is contained in:
+82
-21
@@ -54,13 +54,6 @@ public class EchartsDataUtils {
|
||||
}
|
||||
/**
|
||||
* 构建ECharts图表数据(带时间补全和特殊值处理)
|
||||
* @param list 原始数据列表
|
||||
* @param timeExtractor 时间字段提取函数
|
||||
* @param dataExtractors 数据提取器Map
|
||||
* @param startTime 开始时间字符串(格式:yyyy-MM-dd HH:mm:ss)
|
||||
* @param endTime 结束时间字符串(格式:yyyy-MM-dd HH:mm:ss)
|
||||
* @param <T> 数据类型泛型
|
||||
* @return 包含xData和yData的Map
|
||||
*/
|
||||
public static <T> Map<String, Object> buildEchartsDataAutoPadding(
|
||||
List<T> list,
|
||||
@@ -91,11 +84,44 @@ public class EchartsDataUtils {
|
||||
.sorted(Comparator.comparing(timeExtractor))
|
||||
.collect(Collectors.toList());
|
||||
|
||||
// 自动检测时间间隔
|
||||
// 自动检测时间间隔(用于数据期间)
|
||||
long timeInterval = detectTimeInterval(sortedList, timeExtractor);
|
||||
|
||||
// 生成完整的时间序列
|
||||
List<Date> fullTimeSeries = generateTimeSeries(startDate, endDate, timeInterval);
|
||||
// 获取数据实际的时间范围
|
||||
Date actualStartTime = timeExtractor.apply(sortedList.get(0));
|
||||
Date actualEndTime = timeExtractor.apply(sortedList.get(sortedList.size() - 1));
|
||||
|
||||
// 计算稀疏间隔(用于数据期间外)
|
||||
long totalTimeRange = endDate.getTime() - startDate.getTime();
|
||||
long sparseInterval = totalTimeRange > 12L * 30 * 24 * 60 * 60 * 1000 ?
|
||||
30L * 24 * 60 * 60 * 1000 : 2L * 24 * 60 * 60 * 1000;
|
||||
|
||||
// 生成三段时间序列
|
||||
List<Date> fullTimeSeries = new ArrayList<>();
|
||||
|
||||
// 1. 开始时间到数据开始时间(稀疏间隔)
|
||||
if (startDate.before(actualStartTime)) {
|
||||
List<Date> beforeSeries = generateTimeSeries(startDate, actualStartTime, sparseInterval);
|
||||
fullTimeSeries.addAll(beforeSeries);
|
||||
}
|
||||
|
||||
// 2. 数据开始时间到数据结束时间(正常间隔)
|
||||
List<Date> dataSeries = generateTimeSeries(actualStartTime, actualEndTime, timeInterval);
|
||||
fullTimeSeries.addAll(dataSeries);
|
||||
|
||||
// 3. 数据结束时间到结束时间(稀疏间隔)
|
||||
if (actualEndTime.before(endDate)) {
|
||||
// 调整actualEndTime的下一个点开始,避免重复
|
||||
Calendar cal = Calendar.getInstance();
|
||||
cal.setTime(actualEndTime);
|
||||
cal.setTimeInMillis(cal.getTimeInMillis() + timeInterval);
|
||||
Date nextAfterActualEnd = cal.getTime();
|
||||
|
||||
if (nextAfterActualEnd.before(endDate) || nextAfterActualEnd.equals(endDate)) {
|
||||
List<Date> afterSeries = generateTimeSeries(nextAfterActualEnd, endDate, sparseInterval);
|
||||
fullTimeSeries.addAll(afterSeries);
|
||||
}
|
||||
}
|
||||
|
||||
// 创建时间到数据的映射(考虑时间精度)
|
||||
Map<Long, T> timeDataMap = sortedList.stream()
|
||||
@@ -105,6 +131,12 @@ public class EchartsDataUtils {
|
||||
(a, b) -> a
|
||||
));
|
||||
|
||||
// 检测整个数据集中是否有真实数据
|
||||
boolean hasRealData = checkHasRealData(sortedList, dataExtractors);
|
||||
|
||||
// 特殊处理:查找percentile95的固定值
|
||||
Object fixedPercentile95Value = findFixedValueForPercentile95(sortedList, dataExtractors);
|
||||
|
||||
// 准备X轴和Y轴数据
|
||||
List<String> xAxisData = new ArrayList<>();
|
||||
Map<String, Object> yData = new LinkedHashMap<>();
|
||||
@@ -113,12 +145,6 @@ public class EchartsDataUtils {
|
||||
dataExtractors.keySet().forEach(name ->
|
||||
yData.put(name, new ArrayList<>()));
|
||||
|
||||
// 特殊处理:查找percentile95的固定值
|
||||
Object fixedPercentile95Value = findFixedValueForPercentile95(sortedList, dataExtractors);
|
||||
|
||||
// 检测整个数据集中是否有真实数据
|
||||
boolean hasRealData = checkHasRealData(sortedList, dataExtractors);
|
||||
|
||||
// 记录当前处理的时间点索引
|
||||
int timeIndex = 0;
|
||||
|
||||
@@ -126,6 +152,9 @@ public class EchartsDataUtils {
|
||||
// X轴数据
|
||||
xAxisData.add(parseDateToStr(time));
|
||||
|
||||
// 判断当前时间点是否在数据实际时间范围内
|
||||
boolean isInDataRange = !time.before(actualStartTime) && !time.after(actualEndTime);
|
||||
|
||||
// Y轴数据
|
||||
Long normalizedTime = normalizeTime(time, timeInterval);
|
||||
T item = timeDataMap.get(normalizedTime);
|
||||
@@ -138,11 +167,18 @@ public class EchartsDataUtils {
|
||||
List<Object> seriesData = (List<Object>) yData.get(name);
|
||||
|
||||
if (item != null) {
|
||||
// 有真实数据
|
||||
Object value = extractor.apply(item);
|
||||
seriesData.add(value != null ? value : getDefaultValue(name, fixedPercentile95Value, timeIndex, hasRealData));
|
||||
} else {
|
||||
// 智能数据补全
|
||||
seriesData.add(getDefaultValue(name, fixedPercentile95Value, timeIndex, hasRealData));
|
||||
if (isInDataRange) {
|
||||
// 在数据时间范围内但该时间点无数据:使用智能补全策略
|
||||
seriesData.add(getDefaultValue(name, fixedPercentile95Value, timeIndex, hasRealData));
|
||||
} else {
|
||||
// 在数据时间范围外(开始时间前或结束时间后):使用空数据补全策略
|
||||
seriesData.add(getEmptyDataDefaultValue(name, timeIndex));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -161,6 +197,23 @@ public class EchartsDataUtils {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取空数据默认值(用于数据时间范围外的点)
|
||||
*/
|
||||
private static Object getEmptyDataDefaultValue(String metricName, int timeIndex) {
|
||||
// deployDevice特殊处理:始终补空字符串
|
||||
if ("deployDevice".equals(metricName)) {
|
||||
return "";
|
||||
}
|
||||
|
||||
// 其他字段:第一个点补0,其他点补null
|
||||
if (timeIndex == 0) {
|
||||
return 0;
|
||||
} else {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 检测整个数据集中是否有真实数据(非null)
|
||||
*/
|
||||
@@ -245,9 +298,17 @@ public class EchartsDataUtils {
|
||||
Date endDate = parseStringToDate(endTime);
|
||||
|
||||
if (startDate != null && endDate != null && !startDate.after(endDate)) {
|
||||
// 使用默认时间间隔生成完整时间序列
|
||||
long defaultInterval = 300000L; // 5分钟
|
||||
List<Date> fullTimeSeries = generateTimeSeries(startDate, endDate, defaultInterval);
|
||||
// 动态计算时间间隔
|
||||
long timeRange = endDate.getTime() - startDate.getTime();
|
||||
long interval;
|
||||
|
||||
if (timeRange > 12L * 30 * 24 * 60 * 60 * 1000) { // 超过12个月
|
||||
interval = 30L * 24 * 60 * 60 * 1000; // 每月1个点
|
||||
} else {
|
||||
interval = 2L * 24 * 60 * 60 * 1000; // 2天1个点
|
||||
}
|
||||
|
||||
List<Date> fullTimeSeries = generateTimeSeries(startDate, endDate, interval);
|
||||
|
||||
// 构建x轴数据
|
||||
List<String> xAxisData = new ArrayList<>();
|
||||
@@ -256,7 +317,7 @@ public class EchartsDataUtils {
|
||||
}
|
||||
result.put("xData", xAxisData);
|
||||
|
||||
// 构建y轴数据(空数据集:第一个点补0,其他点补null)
|
||||
// 构建y轴数据(空数据集)
|
||||
Map<String, Object> yData = new LinkedHashMap<>();
|
||||
int dataSize = xAxisData.size();
|
||||
dataNames.forEach(name -> {
|
||||
|
||||
+2
@@ -15,6 +15,8 @@ public class PolicyTypeVo {
|
||||
private String versions;
|
||||
/** 路由信息 */
|
||||
private String routes;
|
||||
/** mtr策略 */
|
||||
private String mtrPolicys;
|
||||
/** 时间戳 */
|
||||
private Long timestamp = Instant.now().getEpochSecond();
|
||||
}
|
||||
|
||||
+2
-1
@@ -5,7 +5,8 @@ import lombok.Getter;
|
||||
@Getter
|
||||
public enum AlarmTypeEnum {
|
||||
服务器下线("1", "服务器下线"),
|
||||
交换机下线("2", "交换机下线");
|
||||
交换机下线("2", "交换机下线"),
|
||||
mtrAgent下线("3", "mtrAgent下线");
|
||||
private final String code;
|
||||
private final String msg;
|
||||
AlarmTypeEnum(String code, String msg){
|
||||
|
||||
+50
-48
@@ -1,24 +1,22 @@
|
||||
package com.ruoyi.mtragent.handler;
|
||||
|
||||
import com.alibaba.fastjson.JSONObject;
|
||||
import com.ruoyi.common.core.constant.SecurityConstants;
|
||||
import com.ruoyi.common.core.domain.R;
|
||||
import com.ruoyi.common.core.enums.MsgEnum;
|
||||
import com.ruoyi.common.core.utils.DateUtils;
|
||||
import com.ruoyi.common.core.utils.StringUtils;
|
||||
import com.ruoyi.mtragent.domain.*;
|
||||
import com.ruoyi.mtragent.domain.vo.MessageVo;
|
||||
import com.ruoyi.mtragent.domain.vo.PolicyTypeVo;
|
||||
import com.ruoyi.mtragent.domain.vo.RegisterMsgVo;
|
||||
import com.ruoyi.mtragent.domain.vo.RspVo;
|
||||
import com.ruoyi.mtragent.enums.AlarmTypeEnum;
|
||||
import com.ruoyi.mtragent.enums.PushMethodEnum;
|
||||
import com.ruoyi.mtragent.model.ProducerMode;
|
||||
import com.ruoyi.mtragent.producer.MessageProducer;
|
||||
import com.ruoyi.mtragent.service.*;
|
||||
import com.ruoyi.mtragent.utils.JsonDataParser;
|
||||
import com.ruoyi.mtragent.utils.WeChatWorkBot;
|
||||
import com.ruoyi.system.api.RemoteRevenueConfigService;
|
||||
import com.ruoyi.system.api.domain.NetworkInfo;
|
||||
import com.ruoyi.system.api.domain.RmSwitchManagementRemote;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.dao.DataAccessException;
|
||||
@@ -57,19 +55,17 @@ public class MessageHandler {
|
||||
@Autowired
|
||||
private RedisTemplate<String, String> redisTemplate;
|
||||
@Autowired
|
||||
private IRmAgentManagementService rmAgentManagementService;
|
||||
@Autowired
|
||||
private RemoteRevenueConfigService remoteRevenueConfigService;
|
||||
@Autowired
|
||||
private IRmAlarmPushConfigService rmAlarmPushConfigService;
|
||||
@Autowired
|
||||
private IRmAlarmLogService rmAlarmLogService;
|
||||
@Autowired
|
||||
private IInitialHeartbeatListenLogService initialHeartbeatListenLogService;
|
||||
@Autowired
|
||||
private IRmNetworkInterfaceService rmNetworkInterfaceService;
|
||||
@Autowired
|
||||
private IRmMtrClientRegistrationService rmMtrClientRegistrationService;
|
||||
@Autowired
|
||||
private IRmMtrPolicyConfigService rmMtrPolicyConfigService;
|
||||
@Autowired
|
||||
private ProducerMode producerMode;
|
||||
|
||||
|
||||
/**
|
||||
@@ -81,6 +77,7 @@ public class MessageHandler {
|
||||
|
||||
// 其他类型消息可以单独注册处理器
|
||||
registerHandler(MsgEnum.注册.getValue(), this::handleRegisterMessage);
|
||||
registerHandler(MsgEnum.获取最新策略.getValue(), this::handleNewPolicyMessage);
|
||||
registerHandler(MsgEnum.心跳上报.getValue(), this::handleHeartbeatMessage);
|
||||
registerHandler(MsgEnum.多公网IP探测.getValue(), this::handleNetWorkDelectMessage);
|
||||
}
|
||||
@@ -177,6 +174,41 @@ public class MessageHandler {
|
||||
}
|
||||
|
||||
// ========== 具体的消息处理方法 ==========
|
||||
/**
|
||||
* 获取最新策略
|
||||
* @param deviceMessage
|
||||
*/
|
||||
private void handleNewPolicyMessage(DeviceMessage deviceMessage) {
|
||||
MessageProducer messageProducer = new MessageProducer();
|
||||
List<RegisterMsgVo> interfaces = JsonDataParser.parseJsonData(deviceMessage.getData(), RegisterMsgVo.class);
|
||||
if(!interfaces.isEmpty()) {
|
||||
RegisterMsgVo registerMsgVo = interfaces.get(0);
|
||||
String mtrClientId = registerMsgVo.getClientId();
|
||||
List<RmMtrPolicyConfig> mtrPolicyConfigList = rmMtrPolicyConfigService.getPoliciesForMtrClient(mtrClientId);
|
||||
if(mtrPolicyConfigList != null && !mtrPolicyConfigList.isEmpty()){
|
||||
// 构建mtrclient消息
|
||||
String mtrPolicyListStr = JSONObject.toJSONString(mtrPolicyConfigList);
|
||||
PolicyTypeVo policyTypeVo = new PolicyTypeVo();
|
||||
policyTypeVo.setMtrPolicys(mtrPolicyListStr);
|
||||
String configJson = JSONObject.toJSONString(policyTypeVo);
|
||||
try {
|
||||
DeviceMessage message = new DeviceMessage();
|
||||
message.setClientId(mtrClientId);
|
||||
message.setData(configJson);
|
||||
message.setDataType(MsgEnum.获取最新策略应答.getValue());
|
||||
|
||||
messageProducer.sendAsyncProducerMessage(
|
||||
producerMode.getAgentTopic(),
|
||||
"",
|
||||
"",
|
||||
JSONObject.toJSONString(message)
|
||||
);
|
||||
} catch (Exception e) {
|
||||
log.error("发送mtr策略失败,mtrClientId: {}", mtrClientId, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 监听心跳
|
||||
@@ -420,56 +452,26 @@ public class MessageHandler {
|
||||
* @param clientId
|
||||
*/
|
||||
private void insertAlarmRecords(String clientId) {
|
||||
// 查询clientId是否为交换机唯一标识
|
||||
boolean isSwitch = false;
|
||||
String switchName = "";
|
||||
RmSwitchManagementRemote rmSwitchManagementRemote = new RmSwitchManagementRemote();
|
||||
rmSwitchManagementRemote.setClientId(clientId);
|
||||
R<List<RmSwitchManagementRemote>> rmSwitchManagementRemoteListR = remoteRevenueConfigService.getSwitchNameByClientId(rmSwitchManagementRemote, SecurityConstants.INNER);
|
||||
if (rmSwitchManagementRemoteListR != null &&
|
||||
rmSwitchManagementRemoteListR.getData() != null &&
|
||||
!rmSwitchManagementRemoteListR.getData().isEmpty()) {
|
||||
isSwitch = true;
|
||||
switchName = rmSwitchManagementRemoteListR.getData().get(0).getSwitchName();
|
||||
}
|
||||
|
||||
// 创建告警日志记录
|
||||
RmAlarmLog rmAlarmLog = createAlarmLog(clientId, isSwitch, switchName);
|
||||
RmAlarmLog rmAlarmLog = createAlarmLog(clientId);
|
||||
|
||||
// 插入告警日志
|
||||
rmAlarmLogService.insertRmAlarmLog(rmAlarmLog);
|
||||
|
||||
// 发送告警推送
|
||||
sendAlarmPush(rmAlarmLog, isSwitch);
|
||||
sendAlarmPush(rmAlarmLog);
|
||||
}
|
||||
|
||||
/**
|
||||
* 创建告警日志记录
|
||||
*/
|
||||
private RmAlarmLog createAlarmLog(String clientId, boolean isSwitch, String switchName) {
|
||||
private RmAlarmLog createAlarmLog(String clientId) {
|
||||
RmAlarmLog rmAlarmLog = new RmAlarmLog();
|
||||
String alarmContent = clientId + "下线";
|
||||
|
||||
if (isSwitch) {
|
||||
rmAlarmLog.setClientId(switchName);
|
||||
rmAlarmLog.setAlarmType("2");
|
||||
} else {
|
||||
// 查询管理网公网ip
|
||||
RmNetworkInterface rmNetworkInterface = new RmNetworkInterface();
|
||||
rmNetworkInterface.setClientId(clientId);
|
||||
rmNetworkInterface.setNewFlag(1);
|
||||
List<RmNetworkInterface> interfaceList = rmNetworkInterfaceService.selectRmNetworkInterfaceList(rmNetworkInterface);
|
||||
if (interfaceList != null && !interfaceList.isEmpty()) {
|
||||
interfaceList.stream()
|
||||
.filter(info -> "2".equals(info.getBindIp()) || "3".equals(info.getBindIp()))
|
||||
.findFirst()
|
||||
.ifPresent(networkInterface -> {
|
||||
rmAlarmLog.setMgmPublicIp(networkInterface.getPublicIp());
|
||||
});
|
||||
}
|
||||
rmAlarmLog.setClientId(clientId);
|
||||
rmAlarmLog.setAlarmType("1");
|
||||
}
|
||||
rmAlarmLog.setClientId(clientId);
|
||||
rmAlarmLog.setAlarmType("3");
|
||||
|
||||
rmAlarmLog.setAlarmContent(alarmContent);
|
||||
rmAlarmLog.setAlarmTime(DateUtils.getNowDate());
|
||||
@@ -479,9 +481,9 @@ public class MessageHandler {
|
||||
/**
|
||||
* 发送告警推送
|
||||
*/
|
||||
private void sendAlarmPush(RmAlarmLog rmAlarmLog, boolean isSwitch) {
|
||||
String alarmTypeCode = isSwitch ? AlarmTypeEnum.交换机下线.getCode() : AlarmTypeEnum.服务器下线.getCode();
|
||||
String alarmTypeMsg = isSwitch ? AlarmTypeEnum.交换机下线.getMsg() : AlarmTypeEnum.服务器下线.getMsg();
|
||||
private void sendAlarmPush(RmAlarmLog rmAlarmLog) {
|
||||
String alarmTypeCode = AlarmTypeEnum.mtrAgent下线.getCode();
|
||||
String alarmTypeMsg = AlarmTypeEnum.mtrAgent下线.getMsg();
|
||||
|
||||
RmAlarmPushConfig rmAlarmPushConfig = new RmAlarmPushConfig();
|
||||
rmAlarmPushConfig.setPushMethod(PushMethodEnum.企业微信.getCode());
|
||||
|
||||
+10
-1
@@ -1,8 +1,9 @@
|
||||
package com.ruoyi.mtragent.service;
|
||||
|
||||
import java.util.List;
|
||||
import com.ruoyi.mtragent.domain.RmMtrPolicyConfig;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* mtr探测策略配置Service接口
|
||||
*
|
||||
@@ -58,4 +59,12 @@ public interface IRmMtrPolicyConfigService
|
||||
* @return 结果
|
||||
*/
|
||||
public int deleteRmMtrPolicyConfigById(Long id);
|
||||
|
||||
/**
|
||||
* 根据MTR客户端ID获取应该下发的策略列表
|
||||
*
|
||||
* @param mtrClientId MTR客户端ID
|
||||
* @return 策略列表
|
||||
*/
|
||||
List<RmMtrPolicyConfig> getPoliciesForMtrClient(String mtrClientId);
|
||||
}
|
||||
|
||||
+155
-2
@@ -12,8 +12,8 @@ import com.ruoyi.system.api.domain.RmNetworkInterfaceRemote;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.StringJoiner;
|
||||
import java.util.*;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
* mtr探测策略配置Service业务层处理
|
||||
@@ -93,6 +93,8 @@ public class RmMtrPolicyConfigServiceImpl implements IRmMtrPolicyConfigService
|
||||
rmMtrPolicyConfig.setCreateTime(DateUtils.getNowDate());
|
||||
rmMtrPolicyConfig.setUpdateTime(DateUtils.getNowDate());
|
||||
rmMtrPolicyConfig.setCreateBy(SecurityUtils.getUsername());
|
||||
// 数据表触发器,优先级自增
|
||||
rmMtrPolicyConfig.setPriority(null);
|
||||
// 给ip赋值
|
||||
setServerip(rmMtrPolicyConfig);
|
||||
int rows = rmMtrPolicyConfigMapper.insertRmMtrPolicyConfig(rmMtrPolicyConfig);
|
||||
@@ -135,4 +137,155 @@ public class RmMtrPolicyConfigServiceImpl implements IRmMtrPolicyConfigService
|
||||
{
|
||||
return rmMtrPolicyConfigMapper.deleteRmMtrPolicyConfigById(id);
|
||||
}
|
||||
/**
|
||||
* 根据MTR客户端ID获取应该下发的策略列表
|
||||
*
|
||||
* @param mtrClientId MTR客户端ID
|
||||
* @return 策略列表
|
||||
*/
|
||||
@Override
|
||||
public List<RmMtrPolicyConfig> getPoliciesForMtrClient(String mtrClientId) {
|
||||
List<RmMtrPolicyConfig> result = new ArrayList<>();
|
||||
if (mtrClientId == null || mtrClientId.trim().isEmpty()) {
|
||||
return result;
|
||||
}
|
||||
|
||||
// 查询所有有效的策略(探测标志为1)
|
||||
RmMtrPolicyConfig query = new RmMtrPolicyConfig();
|
||||
query.setProbeFlag(1L); // 只查询启用探测的策略
|
||||
List<RmMtrPolicyConfig> allPolicies = rmMtrPolicyConfigMapper.selectRmMtrPolicyConfigList(query);
|
||||
|
||||
if (allPolicies == null || allPolicies.isEmpty()) {
|
||||
return result;
|
||||
}
|
||||
|
||||
// 按优先级降序排序(优先级高的在前面)
|
||||
List<RmMtrPolicyConfig> sortedPolicies = allPolicies.stream()
|
||||
.filter(p -> p.getPriority() != null && p.getProbeFlag() != null)
|
||||
.sorted((p1, p2) -> Long.compare(p2.getPriority(), p1.getPriority()))
|
||||
.collect(Collectors.toList());
|
||||
|
||||
// 构建最终的服务器分配映射
|
||||
Map<String, String> serverAssignment = new HashMap<>();
|
||||
|
||||
// 从高优先级到低优先级处理策略
|
||||
for (RmMtrPolicyConfig policy : sortedPolicies) {
|
||||
List<String> serverIds = getServerIdsFromPolicy(policy);
|
||||
if (serverIds.isEmpty()) {
|
||||
continue;
|
||||
}
|
||||
|
||||
// 处理每个服务器:只有这个服务器还没有被处理过时,才处理当前策略
|
||||
for (String serverId : serverIds) {
|
||||
if (!serverAssignment.containsKey(serverId)) {
|
||||
if (policy.getProbeFlag() == 1) {
|
||||
// 探测:记录分配给哪个MTR客户端
|
||||
serverAssignment.put(serverId, policy.getMtrClientId());
|
||||
} else {
|
||||
// 不探测:记录为null
|
||||
serverAssignment.put(serverId, null);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 找出应该由当前请求客户端执行的服务器
|
||||
List<String> assignedServers = new ArrayList<>();
|
||||
for (Map.Entry<String, String> entry : serverAssignment.entrySet()) {
|
||||
String serverId = entry.getKey();
|
||||
String assignedClientId = entry.getValue();
|
||||
|
||||
// 只有当分配给的客户端是当前请求客户端,且不是null(即需要探测)时
|
||||
if (mtrClientId.equals(assignedClientId)) {
|
||||
assignedServers.add(serverId);
|
||||
}
|
||||
}
|
||||
|
||||
// 如果没有服务器分配给当前客户端,返回空列表
|
||||
if (assignedServers.isEmpty()) {
|
||||
return result;
|
||||
}
|
||||
|
||||
// 找出每个服务器最终生效的策略(最高优先级的那个)
|
||||
Map<RmMtrPolicyConfig, List<String>> policyServersMap = new HashMap<>();
|
||||
for (String serverId : assignedServers) {
|
||||
RmMtrPolicyConfig finalPolicy = findFinalPolicyForServer(serverId, sortedPolicies);
|
||||
if (finalPolicy != null) {
|
||||
policyServersMap.computeIfAbsent(finalPolicy, k -> new ArrayList<>()).add(serverId);
|
||||
}
|
||||
}
|
||||
|
||||
// 创建策略副本返回
|
||||
for (Map.Entry<RmMtrPolicyConfig, List<String>> entry : policyServersMap.entrySet()) {
|
||||
RmMtrPolicyConfig newPolicy = createPolicyCopy(entry.getKey(), entry.getValue());
|
||||
setServerip(newPolicy);
|
||||
result.add(newPolicy);
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
/**
|
||||
* 找出服务器最终生效的策略(最高优先级的那个)
|
||||
*/
|
||||
private RmMtrPolicyConfig findFinalPolicyForServer(String serverId, List<RmMtrPolicyConfig> sortedPolicies) {
|
||||
for (RmMtrPolicyConfig policy : sortedPolicies) {
|
||||
if (containsServer(policy, serverId)) {
|
||||
return policy;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* 从策略中提取服务器ID列表
|
||||
*/
|
||||
private List<String> getServerIdsFromPolicy(RmMtrPolicyConfig policy) {
|
||||
List<String> serverIds = new ArrayList<>();
|
||||
if (policy.getServerGroup() != null) {
|
||||
// 先还原换行符分割
|
||||
String serverGroup = policy.getServerGroup().replace(";", "\n");
|
||||
String[] ids = serverGroup.split("\\r?\\n");
|
||||
for (String id : ids) {
|
||||
if (id != null && !id.trim().isEmpty()) {
|
||||
serverIds.add(id.trim());
|
||||
}
|
||||
}
|
||||
}
|
||||
return serverIds;
|
||||
}
|
||||
|
||||
/**
|
||||
* 检查策略是否包含指定服务器
|
||||
*/
|
||||
private boolean containsServer(RmMtrPolicyConfig policy, String serverId) {
|
||||
List<String> serverIds = getServerIdsFromPolicy(policy);
|
||||
return serverIds.contains(serverId);
|
||||
}
|
||||
|
||||
/**
|
||||
* 创建策略副本
|
||||
*/
|
||||
private RmMtrPolicyConfig createPolicyCopy(RmMtrPolicyConfig original, List<String> serverIds) {
|
||||
RmMtrPolicyConfig copy = new RmMtrPolicyConfig();
|
||||
copy.setId(original.getId());
|
||||
copy.setPolicyName(original.getPolicyName());
|
||||
copy.setPriority(original.getPriority());
|
||||
copy.setMtrClientId(original.getMtrClientId());
|
||||
copy.setProbeFlag(original.getProbeFlag());
|
||||
copy.setStartTime(original.getStartTime());
|
||||
copy.setEndTime(original.getEndTime());
|
||||
copy.setProbeFrequency(original.getProbeFrequency());
|
||||
|
||||
// 服务器组用换行符连接
|
||||
copy.setServerGroup(String.join("\n", serverIds));
|
||||
copy.setServeripGroup(original.getServeripGroup());
|
||||
|
||||
copy.setCreateTime(original.getCreateTime());
|
||||
copy.setUpdateTime(original.getUpdateTime());
|
||||
copy.setCreateBy(original.getCreateBy());
|
||||
copy.setUpdateBy(original.getUpdateBy());
|
||||
|
||||
return copy;
|
||||
}
|
||||
}
|
||||
|
||||
+24
-2
@@ -202,7 +202,7 @@ public class RmMonitorConfigServiceImpl implements IRmMonitorConfigService
|
||||
}
|
||||
// 批量插入
|
||||
if(!batchInsertList.isEmpty()){
|
||||
rmMonitorConfigDetailsMapper.batchInsertRmMonitorConfigDetails(batchInsertList);
|
||||
batchInsertInChunks(batchInsertList, rmMonitorConfig.getId());
|
||||
}
|
||||
}else{
|
||||
List<InitialSwitchInfoDetails> switchInfoDetailsList = initialSwitchInfoDetailsService.getMonitorViewDetails(rmMonitorConfig);
|
||||
@@ -229,11 +229,33 @@ public class RmMonitorConfigServiceImpl implements IRmMonitorConfigService
|
||||
}
|
||||
// 批量插入
|
||||
if(!batchInsertList.isEmpty()){
|
||||
rmMonitorConfigDetailsMapper.batchInsertRmMonitorConfigDetails(batchInsertList);
|
||||
batchInsertInChunks(batchInsertList, rmMonitorConfig.getId());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
private void batchInsertInChunks(List<RmMonitorConfigDetails> batchInsertList, Long monitorId) {
|
||||
if (batchInsertList == null || batchInsertList.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
|
||||
int batchSize = 288; // 可配置化
|
||||
int totalSize = batchInsertList.size();
|
||||
|
||||
for (int i = 0; i < totalSize; i += batchSize) {
|
||||
int end = Math.min(i + batchSize, totalSize);
|
||||
List<RmMonitorConfigDetails> subList = batchInsertList.subList(i, end);
|
||||
|
||||
try {
|
||||
rmMonitorConfigDetailsMapper.batchInsertRmMonitorConfigDetails(subList);
|
||||
log.debug("监控配置 {}: 成功插入 {} 条记录,进度: {}/{}",
|
||||
monitorId, subList.size(), i/batchSize + 1,
|
||||
(totalSize + batchSize - 1) / batchSize);
|
||||
} catch (Exception e) {
|
||||
log.error("监控配置 {} 批次插入失败,范围: {}-{}", monitorId, i, end, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
+2
@@ -31,6 +31,8 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
|
||||
<if test="outBytes != null "> and out_bytes = #{outBytes}</if>
|
||||
<if test="inSpeed != null "> and in_speed = #{inSpeed}</if>
|
||||
<if test="outSpeed != null "> and out_speed = #{outSpeed}</if>
|
||||
<if test="startTime != null and startTime != ''"> and create_time >= #{startTime}</if>
|
||||
<if test="endTime != null and endTime != ''"> and create_time <= #{endTime}</if>
|
||||
</where>
|
||||
</select>
|
||||
|
||||
|
||||
Reference in New Issue
Block a user