健康分调整为健康状态

This commit is contained in:
zengqiao
2022-10-29 13:43:33 +08:00
committed by EricZeng
parent e98ec562a2
commit b101cec6fa
30 changed files with 1091 additions and 171 deletions

View File

@@ -5,9 +5,7 @@ import com.didiglobal.logi.log.LogFactory;
import com.xiaojukeji.know.streaming.km.biz.cluster.ClusterZookeepersManager; import com.xiaojukeji.know.streaming.km.biz.cluster.ClusterZookeepersManager;
import com.xiaojukeji.know.streaming.km.common.bean.dto.cluster.ClusterZookeepersOverviewDTO; import com.xiaojukeji.know.streaming.km.common.bean.dto.cluster.ClusterZookeepersOverviewDTO;
import com.xiaojukeji.know.streaming.km.common.bean.entity.cluster.ClusterPhy; import com.xiaojukeji.know.streaming.km.common.bean.entity.cluster.ClusterPhy;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.ZKConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.metrics.ZookeeperMetrics; import com.xiaojukeji.know.streaming.km.common.bean.entity.metrics.ZookeeperMetrics;
import com.xiaojukeji.know.streaming.km.common.bean.entity.param.metric.ZookeeperMetricParam;
import com.xiaojukeji.know.streaming.km.common.bean.entity.result.PaginationResult; import com.xiaojukeji.know.streaming.km.common.bean.entity.result.PaginationResult;
import com.xiaojukeji.know.streaming.km.common.bean.entity.result.Result; import com.xiaojukeji.know.streaming.km.common.bean.entity.result.Result;
import com.xiaojukeji.know.streaming.km.common.bean.entity.result.ResultStatus; import com.xiaojukeji.know.streaming.km.common.bean.entity.result.ResultStatus;
@@ -20,7 +18,6 @@ import com.xiaojukeji.know.streaming.km.common.constant.MsgConstant;
import com.xiaojukeji.know.streaming.km.common.enums.zookeeper.ZKRoleEnum; import com.xiaojukeji.know.streaming.km.common.enums.zookeeper.ZKRoleEnum;
import com.xiaojukeji.know.streaming.km.common.utils.ConvertUtil; import com.xiaojukeji.know.streaming.km.common.utils.ConvertUtil;
import com.xiaojukeji.know.streaming.km.common.utils.PaginationUtil; import com.xiaojukeji.know.streaming.km.common.utils.PaginationUtil;
import com.xiaojukeji.know.streaming.km.common.utils.Tuple;
import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterPhyService; import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterPhyService;
import com.xiaojukeji.know.streaming.km.core.service.version.metrics.ZookeeperMetricVersionItems; import com.xiaojukeji.know.streaming.km.core.service.version.metrics.ZookeeperMetricVersionItems;
import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZnodeService; import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZnodeService;
@@ -30,7 +27,6 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.util.Arrays; import java.util.Arrays;
import java.util.List; import java.util.List;
import java.util.stream.Collectors;
@Service @Service
@@ -56,11 +52,6 @@ public class ClusterZookeepersManagerImpl implements ClusterZookeepersManager {
return Result.buildFromRSAndMsg(ResultStatus.CLUSTER_NOT_EXIST, MsgConstant.getClusterPhyNotExist(clusterPhyId)); return Result.buildFromRSAndMsg(ResultStatus.CLUSTER_NOT_EXIST, MsgConstant.getClusterPhyNotExist(clusterPhyId));
} }
// // TODO
// private Integer healthState;
// private Integer healthCheckPassed;
// private Integer healthCheckTotal;
List<ZookeeperInfo> infoList = zookeeperService.listFromDBByCluster(clusterPhyId); List<ZookeeperInfo> infoList = zookeeperService.listFromDBByCluster(clusterPhyId);
ClusterZookeepersStateVO vo = new ClusterZookeepersStateVO(); ClusterZookeepersStateVO vo = new ClusterZookeepersStateVO();
@@ -90,12 +81,17 @@ public class ClusterZookeepersManagerImpl implements ClusterZookeepersManager {
} }
} }
Result<ZookeeperMetrics> metricsResult = zookeeperMetricService.collectMetricsFromZookeeper(new ZookeeperMetricParam( // 指标获取
Result<ZookeeperMetrics> metricsResult = zookeeperMetricService.batchCollectMetricsFromZookeeper(
clusterPhyId, clusterPhyId,
infoList.stream().filter(elem -> elem.alive()).map(item -> new Tuple<String, Integer>(item.getHost(), item.getPort())).collect(Collectors.toList()), Arrays.asList(
ConvertUtil.str2ObjByJson(clusterPhy.getZkProperties(), ZKConfig.class), ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_WATCH_COUNT,
ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_WATCH_COUNT ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_HEALTH_STATE,
)); ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_HEALTH_CHECK_PASSED,
ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_HEALTH_CHECK_TOTAL
)
);
if (metricsResult.failed()) { if (metricsResult.failed()) {
LOGGER.error( LOGGER.error(
"class=ClusterZookeepersManagerImpl||method=getClusterPhyZookeepersState||clusterPhyId={}||errMsg={}", "class=ClusterZookeepersManagerImpl||method=getClusterPhyZookeepersState||clusterPhyId={}||errMsg={}",
@@ -103,8 +99,12 @@ public class ClusterZookeepersManagerImpl implements ClusterZookeepersManager {
); );
return Result.buildSuc(vo); return Result.buildSuc(vo);
} }
Float watchCount = metricsResult.getData().getMetric(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_WATCH_COUNT);
vo.setWatchCount(watchCount != null? watchCount.intValue(): null); ZookeeperMetrics metrics = metricsResult.getData();
vo.setWatchCount(ConvertUtil.float2Integer(metrics.getMetrics().get(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_WATCH_COUNT)));
vo.setHealthState(ConvertUtil.float2Integer(metrics.getMetrics().get(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_HEALTH_STATE)));
vo.setHealthCheckPassed(ConvertUtil.float2Integer(metrics.getMetrics().get(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_HEALTH_CHECK_PASSED)));
vo.setHealthCheckTotal(ConvertUtil.float2Integer(metrics.getMetrics().get(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_HEALTH_CHECK_TOTAL)));
return Result.buildSuc(vo); return Result.buildSuc(vo);
} }

View File

@@ -47,7 +47,7 @@ public class VersionControlManagerImpl implements VersionControlManager {
@PostConstruct @PostConstruct
public void init(){ public void init(){
defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_HEALTH_SCORE, true)); defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_HEALTH_STATE, true));
defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_FAILED_FETCH_REQ, true)); defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_FAILED_FETCH_REQ, true));
defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_FAILED_PRODUCE_REQ, true)); defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_FAILED_PRODUCE_REQ, true));
defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_UNDER_REPLICA_PARTITIONS, true)); defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_UNDER_REPLICA_PARTITIONS, true));
@@ -57,7 +57,7 @@ public class VersionControlManagerImpl implements VersionControlManager {
defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_BYTES_REJECTED, true)); defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_BYTES_REJECTED, true));
defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_MESSAGE_IN, true)); defaultMetrics.add(new UserMetricConfig(METRIC_TOPIC.getCode(), TOPIC_METRIC_MESSAGE_IN, true));
defaultMetrics.add(new UserMetricConfig(METRIC_CLUSTER.getCode(), CLUSTER_METRIC_HEALTH_SCORE, true)); defaultMetrics.add(new UserMetricConfig(METRIC_CLUSTER.getCode(), CLUSTER_METRIC_HEALTH_STATE, true));
defaultMetrics.add(new UserMetricConfig(METRIC_CLUSTER.getCode(), CLUSTER_METRIC_ACTIVE_CONTROLLER_COUNT, true)); defaultMetrics.add(new UserMetricConfig(METRIC_CLUSTER.getCode(), CLUSTER_METRIC_ACTIVE_CONTROLLER_COUNT, true));
defaultMetrics.add(new UserMetricConfig(METRIC_CLUSTER.getCode(), CLUSTER_METRIC_BYTES_IN, true)); defaultMetrics.add(new UserMetricConfig(METRIC_CLUSTER.getCode(), CLUSTER_METRIC_BYTES_IN, true));
defaultMetrics.add(new UserMetricConfig(METRIC_CLUSTER.getCode(), CLUSTER_METRIC_BYTES_OUT, true)); defaultMetrics.add(new UserMetricConfig(METRIC_CLUSTER.getCode(), CLUSTER_METRIC_BYTES_OUT, true));
@@ -75,9 +75,9 @@ public class VersionControlManagerImpl implements VersionControlManager {
defaultMetrics.add(new UserMetricConfig(METRIC_GROUP.getCode(), GROUP_METRIC_OFFSET_CONSUMED, true)); defaultMetrics.add(new UserMetricConfig(METRIC_GROUP.getCode(), GROUP_METRIC_OFFSET_CONSUMED, true));
defaultMetrics.add(new UserMetricConfig(METRIC_GROUP.getCode(), GROUP_METRIC_LAG, true)); defaultMetrics.add(new UserMetricConfig(METRIC_GROUP.getCode(), GROUP_METRIC_LAG, true));
defaultMetrics.add(new UserMetricConfig(METRIC_GROUP.getCode(), GROUP_METRIC_STATE, true)); defaultMetrics.add(new UserMetricConfig(METRIC_GROUP.getCode(), GROUP_METRIC_STATE, true));
defaultMetrics.add(new UserMetricConfig(METRIC_GROUP.getCode(), GROUP_METRIC_HEALTH_SCORE, true)); defaultMetrics.add(new UserMetricConfig(METRIC_GROUP.getCode(), GROUP_METRIC_HEALTH_STATE, true));
defaultMetrics.add(new UserMetricConfig(METRIC_BROKER.getCode(), BROKER_METRIC_HEALTH_SCORE, true)); defaultMetrics.add(new UserMetricConfig(METRIC_BROKER.getCode(), BROKER_METRIC_HEALTH_STATE, true));
defaultMetrics.add(new UserMetricConfig(METRIC_BROKER.getCode(), BROKER_METRIC_CONNECTION_COUNT, true)); defaultMetrics.add(new UserMetricConfig(METRIC_BROKER.getCode(), BROKER_METRIC_CONNECTION_COUNT, true));
defaultMetrics.add(new UserMetricConfig(METRIC_BROKER.getCode(), BROKER_METRIC_MESSAGE_IN, true)); defaultMetrics.add(new UserMetricConfig(METRIC_BROKER.getCode(), BROKER_METRIC_MESSAGE_IN, true));
defaultMetrics.add(new UserMetricConfig(METRIC_BROKER.getCode(), BROKER_METRIC_NETWORK_RPO_AVG_IDLE, true)); defaultMetrics.add(new UserMetricConfig(METRIC_BROKER.getCode(), BROKER_METRIC_NETWORK_RPO_AVG_IDLE, true));

View File

@@ -3,6 +3,7 @@ package com.xiaojukeji.know.streaming.km.common.bean.entity.broker;
import com.alibaba.fastjson.TypeReference; import com.alibaba.fastjson.TypeReference;
import com.xiaojukeji.know.streaming.km.common.bean.entity.common.IpPortData; import com.xiaojukeji.know.streaming.km.common.bean.entity.common.IpPortData;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.JmxConfig;
import com.xiaojukeji.know.streaming.km.common.bean.po.broker.BrokerPO; import com.xiaojukeji.know.streaming.km.common.bean.po.broker.BrokerPO;
import com.xiaojukeji.know.streaming.km.common.utils.ConvertUtil; import com.xiaojukeji.know.streaming.km.common.utils.ConvertUtil;
import lombok.AllArgsConstructor; import lombok.AllArgsConstructor;
@@ -65,13 +66,13 @@ public class Broker implements Serializable {
*/ */
private Map<String, IpPortData> endpointMap; private Map<String, IpPortData> endpointMap;
public static Broker buildFrom(Long clusterPhyId, Node node, Long startTimestamp) { public static Broker buildFrom(Long clusterPhyId, Node node, Long startTimestamp, JmxConfig jmxConfig) {
Broker metadata = new Broker(); Broker metadata = new Broker();
metadata.setClusterPhyId(clusterPhyId); metadata.setClusterPhyId(clusterPhyId);
metadata.setBrokerId(node.id()); metadata.setBrokerId(node.id());
metadata.setHost(node.host()); metadata.setHost(node.host());
metadata.setPort(node.port()); metadata.setPort(node.port());
metadata.setJmxPort(-1); metadata.setJmxPort(jmxConfig != null ? jmxConfig.getJmxPort() : -1);
metadata.setStartTimestamp(startTimestamp); metadata.setStartTimestamp(startTimestamp);
metadata.setRack(node.rack()); metadata.setRack(node.rack());
metadata.setStatus(1); metadata.setStatus(1);

View File

@@ -13,9 +13,4 @@ public class BaseClusterHealthConfig extends BaseClusterConfigValue {
* 健康检查名称 * 健康检查名称
*/ */
protected HealthCheckNameEnum checkNameEnum; protected HealthCheckNameEnum checkNameEnum;
/**
* 权重
*/
protected Float weight;
} }

View File

@@ -0,0 +1,19 @@
package com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck;
import lombok.Data;
/**
* @author wyb
* @date 2022/10/26
*/
@Data
public class HealthAmountRatioConfig extends BaseClusterHealthConfig {
/**
* 总数
*/
private Integer amount;
/**
* 比例
*/
private Double ratio;
}

View File

@@ -0,0 +1,83 @@
package com.xiaojukeji.know.streaming.km.common.bean.entity.health;
import com.xiaojukeji.know.streaming.km.common.bean.po.health.HealthCheckResultPO;
import com.xiaojukeji.know.streaming.km.common.enums.health.HealthCheckNameEnum;
import com.xiaojukeji.know.streaming.km.common.utils.ValidateUtils;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.stream.Collectors;
@Data
@NoArgsConstructor
public class HealthCheckAggResult {
private HealthCheckNameEnum checkNameEnum;
private List<HealthCheckResultPO> poList;
private Boolean passed;
public HealthCheckAggResult(HealthCheckNameEnum checkNameEnum, List<HealthCheckResultPO> poList) {
this.checkNameEnum = checkNameEnum;
this.poList = poList;
if (!ValidateUtils.isEmptyList(poList) && poList.stream().filter(elem -> elem.getPassed() <= 0).count() <= 0) {
passed = true;
} else {
passed = false;
}
}
public Integer getTotalCount() {
if (poList == null) {
return 0;
}
return poList.size();
}
public Integer getPassedCount() {
if (poList == null) {
return 0;
}
return (int) (poList.stream().filter(elem -> elem.getPassed() > 0).count());
}
/**
* 计算当前检查的健康分
* 比如计算集群Broker健康检查中的某一项的健康分
*/
public Integer calRawHealthScore() {
if (poList == null || poList.isEmpty()) {
return 100;
}
return 100 * this.getPassedCount() / this.getTotalCount();
}
public List<String> getNotPassedResNameList() {
if (poList == null) {
return new ArrayList<>();
}
return poList.stream().filter(elem -> elem.getPassed() <= 0).map(elem -> elem.getResName()).collect(Collectors.toList());
}
public Date getCreateTime() {
if (ValidateUtils.isEmptyList(poList)) {
return null;
}
return poList.get(0).getCreateTime();
}
public Date getUpdateTime() {
if (ValidateUtils.isEmptyList(poList)) {
return null;
}
return poList.get(0).getUpdateTime();
}
}

View File

@@ -17,10 +17,6 @@ import java.util.stream.Collectors;
public class HealthScoreResult { public class HealthScoreResult {
private HealthCheckNameEnum checkNameEnum; private HealthCheckNameEnum checkNameEnum;
private Float presentDimensionTotalWeight;
private Float allDimensionTotalWeight;
private BaseClusterHealthConfig baseConfig; private BaseClusterHealthConfig baseConfig;
private List<HealthCheckResultPO> poList; private List<HealthCheckResultPO> poList;
@@ -28,15 +24,11 @@ public class HealthScoreResult {
private Boolean passed; private Boolean passed;
public HealthScoreResult(HealthCheckNameEnum checkNameEnum, public HealthScoreResult(HealthCheckNameEnum checkNameEnum,
Float presentDimensionTotalWeight,
Float allDimensionTotalWeight,
BaseClusterHealthConfig baseConfig, BaseClusterHealthConfig baseConfig,
List<HealthCheckResultPO> poList) { List<HealthCheckResultPO> poList) {
this.checkNameEnum = checkNameEnum; this.checkNameEnum = checkNameEnum;
this.baseConfig = baseConfig; this.baseConfig = baseConfig;
this.poList = poList; this.poList = poList;
this.presentDimensionTotalWeight = presentDimensionTotalWeight;
this.allDimensionTotalWeight = allDimensionTotalWeight;
if (!ValidateUtils.isEmptyList(poList) && poList.stream().filter(elem -> elem.getPassed() <= 0).count() <= 0) { if (!ValidateUtils.isEmptyList(poList) && poList.stream().filter(elem -> elem.getPassed() <= 0).count() <= 0) {
passed = true; passed = true;
} else { } else {
@@ -59,32 +51,6 @@ public class HealthScoreResult {
return (int) (poList.stream().filter(elem -> elem.getPassed() > 0).count()); return (int) (poList.stream().filter(elem -> elem.getPassed() > 0).count());
} }
/**
* 计算所有检查结果的健康分
* 比如:计算集群健康分
*/
public Float calAllWeightHealthScore() {
Float healthScore = 100 * baseConfig.getWeight() / allDimensionTotalWeight;
if (poList == null || poList.isEmpty()) {
return 0.0f;
}
return healthScore * this.getPassedCount() / this.getTotalCount();
}
/**
* 计算当前维度的健康分
* 比如计算集群Broker健康分
*/
public Float calDimensionWeightHealthScore() {
Float healthScore = 100 * baseConfig.getWeight() / presentDimensionTotalWeight;
if (poList == null || poList.isEmpty()) {
return 0.0f;
}
return healthScore * this.getPassedCount() / this.getTotalCount();
}
/** /**
* 计算当前检查的健康分 * 计算当前检查的健康分
* 比如计算集群Broker健康检查中的某一项的健康分 * 比如计算集群Broker健康检查中的某一项的健康分
@@ -102,7 +68,7 @@ public class HealthScoreResult {
return new ArrayList<>(); return new ArrayList<>();
} }
return poList.stream().filter(elem -> elem.getPassed() <= 0).map(elem -> elem.getResName()).collect(Collectors.toList()); return poList.stream().filter(elem -> elem.getPassed() <= 0 && !ValidateUtils.isBlank(elem.getResName())).map(elem -> elem.getResName()).collect(Collectors.toList());
} }
public Date getCreateTime() { public Date getCreateTime() {

View File

@@ -0,0 +1,26 @@
package com.xiaojukeji.know.streaming.km.common.bean.entity.param.zookeeper;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.ZKConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.param.cluster.ClusterPhyParam;
import com.xiaojukeji.know.streaming.km.common.utils.Tuple;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.List;
/**
* @author didi
*/
@Data
@NoArgsConstructor
public class ZookeeperParam extends ClusterPhyParam {
private List<Tuple<String, Integer>> zkAddressList;
private ZKConfig zkConfig;
public ZookeeperParam(Long clusterPhyId, List<Tuple<String, Integer>> zkAddressList, ZKConfig zkConfig) {
super(clusterPhyId);
this.zkAddressList = zkAddressList;
this.zkConfig = zkConfig;
}
}

View File

@@ -32,9 +32,6 @@ public class HealthCheckConfigVO {
@ApiModelProperty(value="检查说明", example = "Group延迟") @ApiModelProperty(value="检查说明", example = "Group延迟")
private String configDesc; private String configDesc;
@ApiModelProperty(value="权重", example = "10")
private Float weight;
@ApiModelProperty(value="检查配置", example = "100") @ApiModelProperty(value="检查配置", example = "100")
private String value; private String value;
} }

View File

@@ -18,6 +18,9 @@ public class HealthScoreBaseResultVO extends BaseTimeVO {
@ApiModelProperty(value="检查维度", example = "1") @ApiModelProperty(value="检查维度", example = "1")
private Integer dimension; private Integer dimension;
@ApiModelProperty(value="检查维度名称", example = "cluster")
private String dimensionName;
@ApiModelProperty(value="检查名称", example = "Group延迟") @ApiModelProperty(value="检查名称", example = "Group延迟")
private String configName; private String configName;
@@ -27,9 +30,6 @@ public class HealthScoreBaseResultVO extends BaseTimeVO {
@ApiModelProperty(value="检查说明", example = "Group延迟") @ApiModelProperty(value="检查说明", example = "Group延迟")
private String configDesc; private String configDesc;
@ApiModelProperty(value="权重百分比[0-100]", example = "10")
private Integer weightPercent;
@ApiModelProperty(value="得分", example = "100") @ApiModelProperty(value="得分", example = "100")
private Integer score; private Integer score;

View File

@@ -35,14 +35,9 @@ public class Constant {
public static final Integer DEFAULT_SESSION_TIMEOUT_UNIT_MS = 15000; public static final Integer DEFAULT_SESSION_TIMEOUT_UNIT_MS = 15000;
public static final Integer DEFAULT_REQUEST_TIMEOUT_UNIT_MS = 5000; public static final Integer DEFAULT_REQUEST_TIMEOUT_UNIT_MS = 5000;
public static final Float MIN_HEALTH_SCORE = 10f;
/** /**
* 指标相关 * 指标相关
*/ */
public static final Integer DEFAULT_CLUSTER_HEALTH_SCORE = 90;
public static final Integer PER_BATCH_MAX_VALUE = 100; public static final Integer PER_BATCH_MAX_VALUE = 100;
public static final String DEFAULT_USER_NAME = "know-streaming-app"; public static final String DEFAULT_USER_NAME = "know-streaming-app";

View File

@@ -15,24 +15,15 @@ public class HealthScoreVOConverter {
private HealthScoreVOConverter() { private HealthScoreVOConverter() {
} }
public static List<HealthScoreResultDetailVO> convert2HealthScoreResultDetailVOList(List<HealthScoreResult> healthScoreResultList, boolean useGlobalWeight) { public static List<HealthScoreResultDetailVO> convert2HealthScoreResultDetailVOList(List<HealthScoreResult> healthScoreResultList) {
Float globalWeightSum = 1f;
if (!healthScoreResultList.isEmpty()) {
globalWeightSum = healthScoreResultList.get(0).getAllDimensionTotalWeight();
}
List<HealthScoreResultDetailVO> voList = new ArrayList<>(); List<HealthScoreResultDetailVO> voList = new ArrayList<>();
for (HealthScoreResult healthScoreResult: healthScoreResultList) { for (HealthScoreResult healthScoreResult: healthScoreResultList) {
HealthScoreResultDetailVO vo = new HealthScoreResultDetailVO(); HealthScoreResultDetailVO vo = new HealthScoreResultDetailVO();
vo.setDimension(healthScoreResult.getCheckNameEnum().getDimensionEnum().getDimension()); vo.setDimension(healthScoreResult.getCheckNameEnum().getDimensionEnum().getDimension());
vo.setDimensionName(healthScoreResult.getCheckNameEnum().getDimensionEnum().getMessage());
vo.setConfigName(healthScoreResult.getCheckNameEnum().getConfigName()); vo.setConfigName(healthScoreResult.getCheckNameEnum().getConfigName());
vo.setConfigItem(healthScoreResult.getCheckNameEnum().getConfigItem()); vo.setConfigItem(healthScoreResult.getCheckNameEnum().getConfigItem());
vo.setConfigDesc(healthScoreResult.getCheckNameEnum().getConfigDesc()); vo.setConfigDesc(healthScoreResult.getCheckNameEnum().getConfigDesc());
if (useGlobalWeight) {
vo.setWeightPercent(healthScoreResult.getBaseConfig().getWeight().intValue() * 100 / globalWeightSum.intValue());
} else {
vo.setWeightPercent(healthScoreResult.getBaseConfig().getWeight().intValue() * 100 / healthScoreResult.getPresentDimensionTotalWeight().intValue());
}
vo.setScore(healthScoreResult.calRawHealthScore()); vo.setScore(healthScoreResult.calRawHealthScore());
if (healthScoreResult.getTotalCount() <= 0) { if (healthScoreResult.getTotalCount() <= 0) {
@@ -57,9 +48,9 @@ public class HealthScoreVOConverter {
for (HealthScoreResult healthScoreResult: healthScoreResultList) { for (HealthScoreResult healthScoreResult: healthScoreResultList) {
HealthScoreBaseResultVO vo = new HealthScoreBaseResultVO(); HealthScoreBaseResultVO vo = new HealthScoreBaseResultVO();
vo.setDimension(healthScoreResult.getCheckNameEnum().getDimensionEnum().getDimension()); vo.setDimension(healthScoreResult.getCheckNameEnum().getDimensionEnum().getDimension());
vo.setDimensionName(healthScoreResult.getCheckNameEnum().getDimensionEnum().getMessage());
vo.setConfigName(healthScoreResult.getCheckNameEnum().getConfigName()); vo.setConfigName(healthScoreResult.getCheckNameEnum().getConfigName());
vo.setConfigDesc(healthScoreResult.getCheckNameEnum().getConfigDesc()); vo.setConfigDesc(healthScoreResult.getCheckNameEnum().getConfigDesc());
vo.setWeightPercent(healthScoreResult.getBaseConfig().getWeight().intValue() * 100 / healthScoreResult.getPresentDimensionTotalWeight().intValue());
vo.setScore(healthScoreResult.calRawHealthScore()); vo.setScore(healthScoreResult.calRawHealthScore());
vo.setPassed(healthScoreResult.getPassedCount().equals(healthScoreResult.getTotalCount())); vo.setPassed(healthScoreResult.getPassedCount().equals(healthScoreResult.getTotalCount()));
vo.setCheckConfig(convert2HealthCheckConfigVO(ConfigGroupEnum.HEALTH.name(), healthScoreResult.getBaseConfig())); vo.setCheckConfig(convert2HealthCheckConfigVO(ConfigGroupEnum.HEALTH.name(), healthScoreResult.getBaseConfig()));
@@ -86,7 +77,6 @@ public class HealthScoreVOConverter {
vo.setConfigName(config.getCheckNameEnum().getConfigName()); vo.setConfigName(config.getCheckNameEnum().getConfigName());
vo.setConfigItem(config.getCheckNameEnum().getConfigItem()); vo.setConfigItem(config.getCheckNameEnum().getConfigItem());
vo.setConfigDesc(config.getCheckNameEnum().getConfigDesc()); vo.setConfigDesc(config.getCheckNameEnum().getConfigDesc());
vo.setWeight(config.getWeight());
vo.setValue(ConvertUtil.obj2Json(config)); vo.setValue(ConvertUtil.obj2Json(config));
return vo; return vo;
} }

View File

@@ -10,13 +10,15 @@ import lombok.Getter;
public enum HealthCheckDimensionEnum { public enum HealthCheckDimensionEnum {
UNKNOWN(-1, "未知"), UNKNOWN(-1, "未知"),
CLUSTER(0, "Cluster维度"), CLUSTER(0, "Cluster"),
BROKER(1, "Broker维度"), BROKER(1, "Broker"),
TOPIC(2, "Topic维度"), TOPIC(2, "Topic"),
GROUP(3, "消费组维度"), GROUP(3, "Group"),
ZOOKEEPER(4, "Zookeeper"),
; ;

View File

@@ -1,6 +1,7 @@
package com.xiaojukeji.know.streaming.km.common.enums.health; package com.xiaojukeji.know.streaming.km.common.enums.health;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.BaseClusterHealthConfig; import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.BaseClusterHealthConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.HealthAmountRatioConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.HealthCompareValueConfig; import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.HealthCompareValueConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.HealthDetectedInLatestMinutesConfig; import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.HealthDetectedInLatestMinutesConfig;
import com.xiaojukeji.know.streaming.km.common.constant.Constant; import com.xiaojukeji.know.streaming.km.common.constant.Constant;
@@ -19,7 +20,8 @@ public enum HealthCheckNameEnum {
"未知", "未知",
Constant.HC_CONFIG_NAME_PREFIX + "UNKNOWN", Constant.HC_CONFIG_NAME_PREFIX + "UNKNOWN",
"未知", "未知",
BaseClusterHealthConfig.class BaseClusterHealthConfig.class,
false
), ),
CLUSTER_NO_CONTROLLER( CLUSTER_NO_CONTROLLER(
@@ -27,7 +29,8 @@ public enum HealthCheckNameEnum {
"Controller", "Controller",
Constant.HC_CONFIG_NAME_PREFIX + "CLUSTER_NO_CONTROLLER", Constant.HC_CONFIG_NAME_PREFIX + "CLUSTER_NO_CONTROLLER",
"集群Controller数正常", "集群Controller数正常",
HealthCompareValueConfig.class HealthCompareValueConfig.class,
true
), ),
BROKER_REQUEST_QUEUE_FULL( BROKER_REQUEST_QUEUE_FULL(
@@ -35,7 +38,8 @@ public enum HealthCheckNameEnum {
"RequestQueueSize", "RequestQueueSize",
Constant.HC_CONFIG_NAME_PREFIX + "BROKER_REQUEST_QUEUE_FULL", Constant.HC_CONFIG_NAME_PREFIX + "BROKER_REQUEST_QUEUE_FULL",
"Broker-RequestQueueSize指标", "Broker-RequestQueueSize指标",
HealthCompareValueConfig.class HealthCompareValueConfig.class,
false
), ),
BROKER_NETWORK_PROCESSOR_AVG_IDLE_TOO_LOW( BROKER_NETWORK_PROCESSOR_AVG_IDLE_TOO_LOW(
@@ -43,7 +47,8 @@ public enum HealthCheckNameEnum {
"NetworkProcessorAvgIdlePercent", "NetworkProcessorAvgIdlePercent",
Constant.HC_CONFIG_NAME_PREFIX + "BROKER_NETWORK_PROCESSOR_AVG_IDLE_TOO_LOW", Constant.HC_CONFIG_NAME_PREFIX + "BROKER_NETWORK_PROCESSOR_AVG_IDLE_TOO_LOW",
"Broker-NetworkProcessorAvgIdlePercent指标", "Broker-NetworkProcessorAvgIdlePercent指标",
HealthCompareValueConfig.class HealthCompareValueConfig.class,
false
), ),
GROUP_RE_BALANCE_TOO_FREQUENTLY( GROUP_RE_BALANCE_TOO_FREQUENTLY(
@@ -51,7 +56,8 @@ public enum HealthCheckNameEnum {
"Group Re-Balance", "Group Re-Balance",
Constant.HC_CONFIG_NAME_PREFIX + "GROUP_RE_BALANCE_TOO_FREQUENTLY", Constant.HC_CONFIG_NAME_PREFIX + "GROUP_RE_BALANCE_TOO_FREQUENTLY",
"Group re-balance频率", "Group re-balance频率",
HealthDetectedInLatestMinutesConfig.class HealthDetectedInLatestMinutesConfig.class,
false
), ),
TOPIC_NO_LEADER( TOPIC_NO_LEADER(
@@ -59,7 +65,8 @@ public enum HealthCheckNameEnum {
"NoLeader", "NoLeader",
Constant.HC_CONFIG_NAME_PREFIX + "TOPIC_NO_LEADER", Constant.HC_CONFIG_NAME_PREFIX + "TOPIC_NO_LEADER",
"Topic 无Leader数", "Topic 无Leader数",
HealthCompareValueConfig.class HealthCompareValueConfig.class,
false
), ),
TOPIC_UNDER_REPLICA_TOO_LONG( TOPIC_UNDER_REPLICA_TOO_LONG(
@@ -67,9 +74,66 @@ public enum HealthCheckNameEnum {
"UnderReplicaTooLong", "UnderReplicaTooLong",
Constant.HC_CONFIG_NAME_PREFIX + "TOPIC_UNDER_REPLICA_TOO_LONG", Constant.HC_CONFIG_NAME_PREFIX + "TOPIC_UNDER_REPLICA_TOO_LONG",
"Topic 未同步持续时间", "Topic 未同步持续时间",
HealthDetectedInLatestMinutesConfig.class HealthDetectedInLatestMinutesConfig.class,
false
), ),
ZK_BRAIN_SPLIT(
HealthCheckDimensionEnum.ZOOKEEPER,
"BrainSplit",
Constant.HC_CONFIG_NAME_PREFIX + "ZK_BRAIN_SPLIT",
"ZK 脑裂",
HealthCompareValueConfig.class,
true
),
ZK_OUTSTANDING_REQUESTS(
HealthCheckDimensionEnum.ZOOKEEPER,
"OutstandingRequests",
Constant.HC_CONFIG_NAME_PREFIX + "ZK_OUTSTANDING_REQUESTS",
"ZK Outstanding 请求堆积数",
HealthAmountRatioConfig.class,
false
),
ZK_WATCH_COUNT(
HealthCheckDimensionEnum.ZOOKEEPER,
"WatchCount",
Constant.HC_CONFIG_NAME_PREFIX + "ZK_WATCH_COUNT",
"ZK WatchCount 数",
HealthAmountRatioConfig.class,
false
),
ZK_ALIVE_CONNECTIONS(
HealthCheckDimensionEnum.ZOOKEEPER,
"AliveConnections",
Constant.HC_CONFIG_NAME_PREFIX + "ZK_ALIVE_CONNECTIONS",
"ZK 连接数",
HealthAmountRatioConfig.class,
false
),
ZK_APPROXIMATE_DATA_SIZE(
HealthCheckDimensionEnum.ZOOKEEPER,
"ApproximateDataSize",
Constant.HC_CONFIG_NAME_PREFIX + "ZK_APPROXIMATE_DATA_SIZE",
"ZK 数据大小(Byte)",
HealthAmountRatioConfig.class,
false
),
ZK_SENT_RATE(
HealthCheckDimensionEnum.ZOOKEEPER,
"SentRate",
Constant.HC_CONFIG_NAME_PREFIX + "ZK_SENT_RATE",
"ZK 发包数",
HealthAmountRatioConfig.class,
false
),
; ;
/** /**
@@ -97,12 +161,18 @@ public enum HealthCheckNameEnum {
*/ */
private final Class configClazz; private final Class configClazz;
HealthCheckNameEnum(HealthCheckDimensionEnum dimensionEnum, String configItem, String configName, String configDesc, Class configClazz) { /**
* 是可用性检查?
*/
private final boolean availableChecker;
HealthCheckNameEnum(HealthCheckDimensionEnum dimensionEnum, String configItem, String configName, String configDesc, Class configClazz, boolean availableChecker) {
this.dimensionEnum = dimensionEnum; this.dimensionEnum = dimensionEnum;
this.configItem = configItem; this.configItem = configItem;
this.configName = configName; this.configName = configName;
this.configDesc = configDesc; this.configDesc = configDesc;
this.configClazz = configClazz; this.configClazz = configClazz;
this.availableChecker = availableChecker;
} }
public static HealthCheckNameEnum getByName(String configName) { public static HealthCheckNameEnum getByName(String configName) {

View File

@@ -16,7 +16,7 @@ public enum HealthStateEnum {
POOR(2, ""), POOR(2, ""),
DEAD(3, "宕机"), DEAD(3, "Down"),
; ;

View File

@@ -28,7 +28,7 @@ import com.xiaojukeji.know.streaming.km.common.utils.ValidateUtils;
import com.xiaojukeji.know.streaming.km.core.cache.CollectedMetricsLocalCache; import com.xiaojukeji.know.streaming.km.core.cache.CollectedMetricsLocalCache;
import com.xiaojukeji.know.streaming.km.core.service.broker.BrokerMetricService; import com.xiaojukeji.know.streaming.km.core.service.broker.BrokerMetricService;
import com.xiaojukeji.know.streaming.km.core.service.broker.BrokerService; import com.xiaojukeji.know.streaming.km.core.service.broker.BrokerService;
import com.xiaojukeji.know.streaming.km.core.service.health.score.HealthScoreService; import com.xiaojukeji.know.streaming.km.core.service.health.state.HealthStateService;
import com.xiaojukeji.know.streaming.km.core.service.partition.PartitionService; import com.xiaojukeji.know.streaming.km.core.service.partition.PartitionService;
import com.xiaojukeji.know.streaming.km.core.service.replica.ReplicaMetricService; import com.xiaojukeji.know.streaming.km.core.service.replica.ReplicaMetricService;
import com.xiaojukeji.know.streaming.km.core.service.topic.TopicService; import com.xiaojukeji.know.streaming.km.core.service.topic.TopicService;
@@ -82,7 +82,7 @@ public class BrokerMetricServiceImpl extends BaseMetricService implements Broker
private ReplicaMetricService replicaMetricService; private ReplicaMetricService replicaMetricService;
@Autowired @Autowired
private HealthScoreService healthScoreService; private HealthStateService healthStateService;
@Autowired @Autowired
private KafkaJMXClient kafkaJMXClient; private KafkaJMXClient kafkaJMXClient;
@@ -108,7 +108,6 @@ public class BrokerMetricServiceImpl extends BaseMetricService implements Broker
registerVCHandler( BROKER_METHOD_GET_HEALTH_SCORE, this::getMetricHealthScore); registerVCHandler( BROKER_METHOD_GET_HEALTH_SCORE, this::getMetricHealthScore);
registerVCHandler( BROKER_METHOD_GET_PARTITIONS_SKEW, this::getPartitionsSkew); registerVCHandler( BROKER_METHOD_GET_PARTITIONS_SKEW, this::getPartitionsSkew);
registerVCHandler( BROKER_METHOD_GET_LEADERS_SKEW, this::getLeadersSkew); registerVCHandler( BROKER_METHOD_GET_LEADERS_SKEW, this::getLeadersSkew);
// registerVCHandler( BROKER_METHOD_GET_LOG_SIZE, this::getLogSize);
registerVCHandler( BROKER_METHOD_GET_LOG_SIZE, V_0_10_0_0, V_1_0_0, "getLogSizeFromJmx", this::getLogSizeFromJmx); registerVCHandler( BROKER_METHOD_GET_LOG_SIZE, V_0_10_0_0, V_1_0_0, "getLogSizeFromJmx", this::getLogSizeFromJmx);
registerVCHandler( BROKER_METHOD_GET_LOG_SIZE, V_1_0_0, V_MAX, "getLogSizeFromClient", this::getLogSizeFromClient); registerVCHandler( BROKER_METHOD_GET_LOG_SIZE, V_1_0_0, V_MAX, "getLogSizeFromClient", this::getLogSizeFromClient);
@@ -318,7 +317,7 @@ public class BrokerMetricServiceImpl extends BaseMetricService implements Broker
Long clusterId = param.getClusterId(); Long clusterId = param.getClusterId();
Integer brokerId = param.getBrokerId(); Integer brokerId = param.getBrokerId();
BrokerMetrics brokerMetrics = healthScoreService.calBrokerHealthScore(clusterId, brokerId); BrokerMetrics brokerMetrics = healthStateService.calBrokerHealthMetrics(clusterId, brokerId);
return Result.buildSuc(brokerMetrics); return Result.buildSuc(brokerMetrics);
} }

View File

@@ -134,7 +134,7 @@ public class BrokerServiceImpl extends BaseVersionControlService implements Brok
newBrokerPO.setId(inDBBrokerPO.getId()); newBrokerPO.setId(inDBBrokerPO.getId());
newBrokerPO.setStatus(Constant.ALIVE); newBrokerPO.setStatus(Constant.ALIVE);
newBrokerPO.setCreateTime(inDBBrokerPO.getCreateTime()); newBrokerPO.setCreateTime(inDBBrokerPO.getCreateTime());
newBrokerPO.setUpdateTime(inDBBrokerPO.getUpdateTime()); newBrokerPO.setUpdateTime(new Date());
if (newBrokerPO.getStartTimestamp() == null) { if (newBrokerPO.getStartTimestamp() == null) {
// 如果当前broker获取不到启动时间 // 如果当前broker获取不到启动时间
// 如果DB中的broker状态为down则使用当前时间否则使用db中已有broker的时间 // 如果DB中的broker状态为down则使用当前时间否则使用db中已有broker的时间
@@ -363,11 +363,11 @@ public class BrokerServiceImpl extends BaseVersionControlService implements Brok
try { try {
Long startTime = jmxDAO.getServerStartTime(clusterPhyId, newNode.host(), null, jmxConfig); Long startTime = jmxDAO.getServerStartTime(clusterPhyId, newNode.host(), null, jmxConfig);
return Broker.buildFrom(clusterPhyId, newNode, startTime); return Broker.buildFrom(clusterPhyId, newNode, startTime, jmxConfig);
} catch (Exception e) { } catch (Exception e) {
log.error("class=BrokerServiceImpl||method=getStartTimeAndBuildBroker||clusterPhyId={}||brokerNode={}||jmxConfig={}||errMsg=exception!", clusterPhyId, newNode, jmxConfig, e); log.error("class=BrokerServiceImpl||method=getStartTimeAndBuildBroker||clusterPhyId={}||brokerNode={}||jmxConfig={}||errMsg=exception!", clusterPhyId, newNode, jmxConfig, e);
} }
return Broker.buildFrom(clusterPhyId, newNode, null); return Broker.buildFrom(clusterPhyId, newNode, null, jmxConfig);
} }
} }

View File

@@ -39,7 +39,7 @@ import com.xiaojukeji.know.streaming.km.core.service.broker.BrokerService;
import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterMetricService; import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterMetricService;
import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterPhyService; import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterPhyService;
import com.xiaojukeji.know.streaming.km.core.service.group.GroupService; import com.xiaojukeji.know.streaming.km.core.service.group.GroupService;
import com.xiaojukeji.know.streaming.km.core.service.health.score.HealthScoreService; import com.xiaojukeji.know.streaming.km.core.service.health.state.HealthStateService;
import com.xiaojukeji.know.streaming.km.core.service.job.JobService; import com.xiaojukeji.know.streaming.km.core.service.job.JobService;
import com.xiaojukeji.know.streaming.km.core.service.kafkacontroller.KafkaControllerService; import com.xiaojukeji.know.streaming.km.core.service.kafkacontroller.KafkaControllerService;
import com.xiaojukeji.know.streaming.km.core.service.partition.PartitionService; import com.xiaojukeji.know.streaming.km.core.service.partition.PartitionService;
@@ -85,7 +85,7 @@ public class ClusterMetricServiceImpl extends BaseMetricService implements Clust
public static final String CLUSTER_METHOD_GET_TOTAL_LOG_SIZE = "getTotalLogSize"; public static final String CLUSTER_METHOD_GET_TOTAL_LOG_SIZE = "getTotalLogSize";
public static final String CLUSTER_METHOD_GET_PARTITION_SIZE = "getPartitionSize"; public static final String CLUSTER_METHOD_GET_PARTITION_SIZE = "getPartitionSize";
public static final String CLUSTER_METHOD_GET_PARTITION_NO_LEADER_SIZE = "getPartitionNoLeaderSize"; public static final String CLUSTER_METHOD_GET_PARTITION_NO_LEADER_SIZE = "getPartitionNoLeaderSize";
public static final String CLUSTER_METHOD_GET_HEALTH_SCORE = "getClusterHealthScore"; public static final String CLUSTER_METHOD_GET_HEALTH_METRICS = "getClusterHealthMetrics";
public static final String CLUSTER_METHOD_GET_METRIC_FROM_KAFKA_BY_TOTAL_BROKERS_JMX = "getMetricFromKafkaByTotalBrokersJMX"; public static final String CLUSTER_METHOD_GET_METRIC_FROM_KAFKA_BY_TOTAL_BROKERS_JMX = "getMetricFromKafkaByTotalBrokersJMX";
public static final String CLUSTER_METHOD_GET_METRIC_FROM_KAFKA_BY_CONTROLLER_JMX = "getMetricFromKafkaByControllerJMX"; public static final String CLUSTER_METHOD_GET_METRIC_FROM_KAFKA_BY_CONTROLLER_JMX = "getMetricFromKafkaByControllerJMX";
public static final String CLUSTER_METHOD_GET_ZK_COUNT = "getZKCount"; public static final String CLUSTER_METHOD_GET_ZK_COUNT = "getZKCount";
@@ -114,7 +114,7 @@ public class ClusterMetricServiceImpl extends BaseMetricService implements Clust
public static final String CLUSTER_METHOD_GET_JOBS_FAILED = "getJobsFailed"; public static final String CLUSTER_METHOD_GET_JOBS_FAILED = "getJobsFailed";
@Autowired @Autowired
private HealthScoreService healthScoreService; private HealthStateService healthStateService;
@Autowired @Autowired
private BrokerService brokerService; private BrokerService brokerService;
@@ -188,7 +188,7 @@ public class ClusterMetricServiceImpl extends BaseMetricService implements Clust
registerVCHandler( CLUSTER_METHOD_GET_PARTITION_SIZE, this::getPartitionSize); registerVCHandler( CLUSTER_METHOD_GET_PARTITION_SIZE, this::getPartitionSize);
registerVCHandler( CLUSTER_METHOD_GET_PARTITION_NO_LEADER_SIZE, this::getPartitionNoLeaderSize); registerVCHandler( CLUSTER_METHOD_GET_PARTITION_NO_LEADER_SIZE, this::getPartitionNoLeaderSize);
registerVCHandler( CLUSTER_METHOD_GET_HEALTH_SCORE, this::getClusterHealthScore); registerVCHandler( CLUSTER_METHOD_GET_HEALTH_METRICS, this::getClusterHealthMetrics);
registerVCHandler( CLUSTER_METHOD_GET_METRIC_FROM_KAFKA_BY_TOTAL_BROKERS_JMX, this::getMetricFromKafkaByTotalBrokersJMX); registerVCHandler( CLUSTER_METHOD_GET_METRIC_FROM_KAFKA_BY_TOTAL_BROKERS_JMX, this::getMetricFromKafkaByTotalBrokersJMX);
registerVCHandler( CLUSTER_METHOD_GET_METRIC_FROM_KAFKA_BY_CONTROLLER_JMX, this::getMetricFromKafkaByControllerJMX); registerVCHandler( CLUSTER_METHOD_GET_METRIC_FROM_KAFKA_BY_CONTROLLER_JMX, this::getMetricFromKafkaByControllerJMX);
@@ -361,15 +361,13 @@ public class ClusterMetricServiceImpl extends BaseMetricService implements Clust
} }
/** private Result<ClusterMetrics> getClusterHealthMetrics(VersionItemParam metricParam){
* 获取集群的健康分
*/
private Result<ClusterMetrics> getClusterHealthScore(VersionItemParam metricParam){
ClusterMetricParam param = (ClusterMetricParam)metricParam; ClusterMetricParam param = (ClusterMetricParam)metricParam;
ClusterMetrics clusterMetrics = healthScoreService.calClusterHealthScore(param.getClusterId()); ClusterMetrics clusterMetrics = healthStateService.calClusterHealthMetrics(param.getClusterId());
return Result.buildSuc(clusterMetrics); return Result.buildSuc(clusterMetrics);
} }
/** /**
* 获取集群的 totalLogSize * 获取集群的 totalLogSize
* @param metricParam * @param metricParam

View File

@@ -20,7 +20,7 @@ import com.xiaojukeji.know.streaming.km.common.utils.BeanUtil;
import com.xiaojukeji.know.streaming.km.common.utils.ConvertUtil; import com.xiaojukeji.know.streaming.km.common.utils.ConvertUtil;
import com.xiaojukeji.know.streaming.km.core.service.group.GroupMetricService; import com.xiaojukeji.know.streaming.km.core.service.group.GroupMetricService;
import com.xiaojukeji.know.streaming.km.core.service.group.GroupService; import com.xiaojukeji.know.streaming.km.core.service.group.GroupService;
import com.xiaojukeji.know.streaming.km.core.service.health.score.HealthScoreService; import com.xiaojukeji.know.streaming.km.core.service.health.state.HealthStateService;
import com.xiaojukeji.know.streaming.km.core.service.partition.PartitionService; import com.xiaojukeji.know.streaming.km.core.service.partition.PartitionService;
import com.xiaojukeji.know.streaming.km.core.service.version.BaseMetricService; import com.xiaojukeji.know.streaming.km.core.service.version.BaseMetricService;
import com.xiaojukeji.know.streaming.km.persistence.es.dao.GroupMetricESDAO; import com.xiaojukeji.know.streaming.km.persistence.es.dao.GroupMetricESDAO;
@@ -64,7 +64,7 @@ public class GroupMetricServiceImpl extends BaseMetricService implements GroupMe
private GroupService groupService; private GroupService groupService;
@Autowired @Autowired
private HealthScoreService healthScoreService; private HealthStateService healthStateService;
@Autowired @Autowired
private PartitionService partitionService; private PartitionService partitionService;
@@ -265,8 +265,8 @@ public class GroupMetricServiceImpl extends BaseMetricService implements GroupMe
private Result<List<GroupMetrics>> getMetricHealthScore(VersionItemParam param) { private Result<List<GroupMetrics>> getMetricHealthScore(VersionItemParam param) {
GroupMetricParam groupMetricParam = (GroupMetricParam)param; GroupMetricParam groupMetricParam = (GroupMetricParam)param;
return Result.buildSuc(Arrays.asList(healthScoreService.calGroupHealthScore( return Result.buildSuc(Arrays.asList(
groupMetricParam.getClusterPhyId(), groupMetricParam.getGroupName())) healthStateService.calGroupHealthMetrics(groupMetricParam.getClusterPhyId(), groupMetricParam.getGroupName()))
); );
} }
} }

View File

@@ -0,0 +1,281 @@
package com.xiaojukeji.know.streaming.km.core.service.health.checker.zookeeper;
import com.didiglobal.logi.log.ILog;
import com.didiglobal.logi.log.LogFactory;
import com.xiaojukeji.know.streaming.km.common.bean.entity.cluster.ClusterPhy;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.ZKConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.BaseClusterHealthConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.HealthAmountRatioConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.HealthCompareValueConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.health.HealthCheckResult;
import com.xiaojukeji.know.streaming.km.common.bean.entity.metrics.ZookeeperMetrics;
import com.xiaojukeji.know.streaming.km.common.bean.entity.param.cluster.ClusterPhyParam;
import com.xiaojukeji.know.streaming.km.common.bean.entity.param.metric.ZookeeperMetricParam;
import com.xiaojukeji.know.streaming.km.common.bean.entity.param.zookeeper.ZookeeperParam;
import com.xiaojukeji.know.streaming.km.common.bean.entity.result.Result;
import com.xiaojukeji.know.streaming.km.common.bean.entity.zookeeper.ZookeeperInfo;
import com.xiaojukeji.know.streaming.km.common.constant.Constant;
import com.xiaojukeji.know.streaming.km.common.enums.health.HealthCheckDimensionEnum;
import com.xiaojukeji.know.streaming.km.common.enums.health.HealthCheckNameEnum;
import com.xiaojukeji.know.streaming.km.common.enums.zookeeper.ZKRoleEnum;
import com.xiaojukeji.know.streaming.km.common.utils.ConvertUtil;
import com.xiaojukeji.know.streaming.km.common.utils.Tuple;
import com.xiaojukeji.know.streaming.km.common.utils.zookeeper.ZookeeperUtils;
import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterPhyService;
import com.xiaojukeji.know.streaming.km.core.service.health.checker.AbstractHealthCheckService;
import com.xiaojukeji.know.streaming.km.core.service.version.metrics.ZookeeperMetricVersionItems;
import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZookeeperMetricService;
import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZookeeperService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import javax.annotation.PostConstruct;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
@Service
public class HealthCheckZookeeperService extends AbstractHealthCheckService {
private static final ILog log = LogFactory.getLog(HealthCheckZookeeperService.class);
@Autowired
private ClusterPhyService clusterPhyService;
@Autowired
private ZookeeperService zookeeperService;
@Autowired
private ZookeeperMetricService zookeeperMetricService;
@PostConstruct
private void init() {
functionMap.putIfAbsent(HealthCheckNameEnum.ZK_BRAIN_SPLIT.getConfigName(), this::checkBrainSplit);
functionMap.putIfAbsent(HealthCheckNameEnum.ZK_OUTSTANDING_REQUESTS.getConfigName(), this::checkOutstandingRequests);
functionMap.putIfAbsent(HealthCheckNameEnum.ZK_WATCH_COUNT.getConfigName(), this::checkWatchCount);
functionMap.putIfAbsent(HealthCheckNameEnum.ZK_ALIVE_CONNECTIONS.getConfigName(), this::checkAliveConnections);
functionMap.putIfAbsent(HealthCheckNameEnum.ZK_APPROXIMATE_DATA_SIZE.getConfigName(), this::checkApproximateDataSize);
functionMap.putIfAbsent(HealthCheckNameEnum.ZK_SENT_RATE.getConfigName(), this::checkSentRate);
}
@Override
public List<ClusterPhyParam> getResList(Long clusterPhyId) {
ClusterPhy clusterPhy = clusterPhyService.getClusterByCluster(clusterPhyId);
if (clusterPhy == null) {
return new ArrayList<>();
}
try {
return Arrays.asList(new ZookeeperParam(
clusterPhyId,
ZookeeperUtils.connectStringParser(clusterPhy.getZookeeper()),
ConvertUtil.str2ObjByJson(clusterPhy.getZkProperties(), ZKConfig.class)
));
} catch (Exception e) {
log.error("class=HealthCheckZookeeperService||method=getResList||clusterPhyId={}||errMsg=exception!", clusterPhyId, e);
}
return new ArrayList<>();
}
@Override
public HealthCheckDimensionEnum getHealthCheckDimensionEnum() {
return HealthCheckDimensionEnum.ZOOKEEPER;
}
private HealthCheckResult checkBrainSplit(Tuple<ClusterPhyParam, BaseClusterHealthConfig> singleConfigSimpleTuple) {
ZookeeperParam param = (ZookeeperParam) singleConfigSimpleTuple.getV1();
HealthCompareValueConfig valueConfig = (HealthCompareValueConfig) singleConfigSimpleTuple.getV2();
List<ZookeeperInfo> infoList = zookeeperService.listFromDBByCluster(param.getClusterPhyId());
HealthCheckResult checkResult = new HealthCheckResult(
HealthCheckDimensionEnum.ZOOKEEPER.getDimension(),
HealthCheckNameEnum.ZK_BRAIN_SPLIT.getConfigName(),
param.getClusterPhyId(),
""
);
long value = infoList.stream().filter(elem -> ZKRoleEnum.LEADER.getRole().equals(elem.getRole())).count();
checkResult.setPassed(value == valueConfig.getValue().longValue() ? Constant.YES : Constant.NO);
return checkResult;
}
private HealthCheckResult checkOutstandingRequests(Tuple<ClusterPhyParam, BaseClusterHealthConfig> singleConfigSimpleTuple) {
ZookeeperParam param = (ZookeeperParam) singleConfigSimpleTuple.getV1();
HealthAmountRatioConfig valueConfig = (HealthAmountRatioConfig) singleConfigSimpleTuple.getV2();
Result<ZookeeperMetrics> metricsResult = zookeeperMetricService.collectMetricsFromZookeeper(
new ZookeeperMetricParam(
param.getClusterPhyId(),
param.getZkAddressList(),
param.getZkConfig(),
ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_OUTSTANDING_REQUESTS
)
);
if (metricsResult.failed() || !metricsResult.hasData()) {
log.error(
"class=HealthCheckZookeeperService||method=checkOutstandingRequests||param={}||config={}||result={}||errMsg=get metrics failed",
param, valueConfig, metricsResult
);
return null;
}
HealthCheckResult checkResult = new HealthCheckResult(
HealthCheckDimensionEnum.ZOOKEEPER.getDimension(),
HealthCheckNameEnum.ZK_OUTSTANDING_REQUESTS.getConfigName(),
param.getClusterPhyId(),
""
);
Float value = metricsResult.getData().getMetric(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_OUTSTANDING_REQUESTS);
checkResult.setPassed(value.intValue() <= valueConfig.getAmount().doubleValue() * valueConfig.getRatio().doubleValue() ? Constant.YES : Constant.NO);
return checkResult;
}
private HealthCheckResult checkWatchCount(Tuple<ClusterPhyParam, BaseClusterHealthConfig> singleConfigSimpleTuple) {
ZookeeperParam param = (ZookeeperParam) singleConfigSimpleTuple.getV1();
HealthAmountRatioConfig valueConfig = (HealthAmountRatioConfig) singleConfigSimpleTuple.getV2();
Result<ZookeeperMetrics> metricsResult = zookeeperMetricService.collectMetricsFromZookeeper(
new ZookeeperMetricParam(
param.getClusterPhyId(),
param.getZkAddressList(),
param.getZkConfig(),
ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_WATCH_COUNT
)
);
if (metricsResult.failed() || !metricsResult.hasData()) {
log.error(
"class=HealthCheckZookeeperService||method=checkWatchCount||param={}||config={}||result={}||errMsg=get metrics failed",
param, valueConfig, metricsResult
);
return null;
}
HealthCheckResult checkResult = new HealthCheckResult(
HealthCheckDimensionEnum.ZOOKEEPER.getDimension(),
HealthCheckNameEnum.ZK_WATCH_COUNT.getConfigName(),
param.getClusterPhyId(),
""
);
Float value = metricsResult.getData().getMetric(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_WATCH_COUNT);
checkResult.setPassed(value.intValue() <= valueConfig.getAmount().doubleValue() * valueConfig.getRatio().doubleValue() ? Constant.YES : Constant.NO);
return checkResult;
}
private HealthCheckResult checkAliveConnections(Tuple<ClusterPhyParam, BaseClusterHealthConfig> singleConfigSimpleTuple) {
ZookeeperParam param = (ZookeeperParam) singleConfigSimpleTuple.getV1();
HealthAmountRatioConfig valueConfig = (HealthAmountRatioConfig) singleConfigSimpleTuple.getV2();
Result<ZookeeperMetrics> metricsResult = zookeeperMetricService.collectMetricsFromZookeeper(
new ZookeeperMetricParam(
param.getClusterPhyId(),
param.getZkAddressList(),
param.getZkConfig(),
ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_NUM_ALIVE_CONNECTIONS
)
);
if (metricsResult.failed() || !metricsResult.hasData()) {
log.error(
"class=HealthCheckZookeeperService||method=checkAliveConnections||param={}||config={}||result={}||errMsg=get metrics failed",
param, valueConfig, metricsResult
);
return null;
}
HealthCheckResult checkResult = new HealthCheckResult(
HealthCheckDimensionEnum.ZOOKEEPER.getDimension(),
HealthCheckNameEnum.ZK_ALIVE_CONNECTIONS.getConfigName(),
param.getClusterPhyId(),
""
);
Float value = metricsResult.getData().getMetric(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_NUM_ALIVE_CONNECTIONS);
checkResult.setPassed(value.intValue() <= valueConfig.getAmount().doubleValue() * valueConfig.getRatio().doubleValue() ? Constant.YES : Constant.NO);
return checkResult;
}
private HealthCheckResult checkApproximateDataSize(Tuple<ClusterPhyParam, BaseClusterHealthConfig> singleConfigSimpleTuple) {
ZookeeperParam param = (ZookeeperParam) singleConfigSimpleTuple.getV1();
HealthAmountRatioConfig valueConfig = (HealthAmountRatioConfig) singleConfigSimpleTuple.getV2();
Result<ZookeeperMetrics> metricsResult = zookeeperMetricService.collectMetricsFromZookeeper(
new ZookeeperMetricParam(
param.getClusterPhyId(),
param.getZkAddressList(),
param.getZkConfig(),
ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_APPROXIMATE_DATA_SIZE
)
);
if (metricsResult.failed() || !metricsResult.hasData()) {
log.error(
"class=HealthCheckZookeeperService||method=checkApproximateDataSize||param={}||config={}||result={}||errMsg=get metrics failed",
param, valueConfig, metricsResult
);
return null;
}
HealthCheckResult checkResult = new HealthCheckResult(
HealthCheckDimensionEnum.ZOOKEEPER.getDimension(),
HealthCheckNameEnum.ZK_APPROXIMATE_DATA_SIZE.getConfigName(),
param.getClusterPhyId(),
""
);
Float value = metricsResult.getData().getMetric(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_APPROXIMATE_DATA_SIZE);
checkResult.setPassed(value.intValue() <= valueConfig.getAmount().doubleValue() * valueConfig.getRatio().doubleValue() ? Constant.YES : Constant.NO);
return checkResult;
}
private HealthCheckResult checkSentRate(Tuple<ClusterPhyParam, BaseClusterHealthConfig> singleConfigSimpleTuple) {
ZookeeperParam param = (ZookeeperParam) singleConfigSimpleTuple.getV1();
HealthAmountRatioConfig valueConfig = (HealthAmountRatioConfig) singleConfigSimpleTuple.getV2();
Result<ZookeeperMetrics> metricsResult = zookeeperMetricService.collectMetricsFromZookeeper(
new ZookeeperMetricParam(
param.getClusterPhyId(),
param.getZkAddressList(),
param.getZkConfig(),
ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_PACKETS_SENT
)
);
if (metricsResult.failed() || !metricsResult.hasData()) {
log.error(
"class=HealthCheckZookeeperService||method=checkSentRate||param={}||config={}||result={}||errMsg=get metrics failed",
param, valueConfig, metricsResult
);
return null;
}
HealthCheckResult checkResult = new HealthCheckResult(
HealthCheckDimensionEnum.ZOOKEEPER.getDimension(),
HealthCheckNameEnum.ZK_SENT_RATE.getConfigName(),
param.getClusterPhyId(),
""
);
Float value = metricsResult.getData().getMetric(ZookeeperMetricVersionItems.ZOOKEEPER_METRIC_PACKETS_SENT);
checkResult.setPassed(value.intValue() <= valueConfig.getAmount().doubleValue() * valueConfig.getRatio().doubleValue() ? Constant.YES : Constant.NO);
return checkResult;
}
}

View File

@@ -0,0 +1,50 @@
package com.xiaojukeji.know.streaming.km.core.service.health.state;
import com.xiaojukeji.know.streaming.km.common.bean.entity.health.HealthScoreResult;
import com.xiaojukeji.know.streaming.km.common.bean.entity.metrics.*;
import com.xiaojukeji.know.streaming.km.common.enums.health.HealthCheckDimensionEnum;
import java.util.List;
public interface HealthStateService {
/**
* 集群健康指标
*/
ClusterMetrics calClusterHealthMetrics(Long clusterPhyId);
/**
* 获取Broker健康指标
*/
BrokerMetrics calBrokerHealthMetrics(Long clusterPhyId, Integer brokerId);
/**
* 获取Topic健康指标
*/
TopicMetrics calTopicHealthMetrics(Long clusterPhyId, String topicName);
/**
* 获取Group健康指标
*/
GroupMetrics calGroupHealthMetrics(Long clusterPhyId, String groupName);
/**
* 获取Zookeeper健康指标
*/
ZookeeperMetrics calZookeeperHealthMetrics(Long clusterPhyId);
/**
* 获取集群健康检查结果
*/
List<HealthScoreResult> getClusterHealthResult(Long clusterPhyId);
/**
* 获取集群某个维度健康检查结果
*/
List<HealthScoreResult> getDimensionHealthResult(Long clusterPhyId, HealthCheckDimensionEnum dimensionEnum);
/**
* 获取集群某个资源的健康检查结果
*/
List<HealthScoreResult> getResHealthResult(Long clusterPhyId, Integer dimension, String resNme);
}

View File

@@ -0,0 +1,408 @@
package com.xiaojukeji.know.streaming.km.core.service.health.state.impl;
import com.xiaojukeji.know.streaming.km.common.bean.entity.broker.Broker;
import com.xiaojukeji.know.streaming.km.common.bean.entity.config.healthcheck.BaseClusterHealthConfig;
import com.xiaojukeji.know.streaming.km.common.bean.entity.health.HealthCheckAggResult;
import com.xiaojukeji.know.streaming.km.common.bean.entity.health.HealthScoreResult;
import com.xiaojukeji.know.streaming.km.common.bean.entity.metrics.*;
import com.xiaojukeji.know.streaming.km.common.bean.po.health.HealthCheckResultPO;
import com.xiaojukeji.know.streaming.km.common.enums.health.HealthCheckDimensionEnum;
import com.xiaojukeji.know.streaming.km.common.enums.health.HealthCheckNameEnum;
import com.xiaojukeji.know.streaming.km.common.enums.health.HealthStateEnum;
import com.xiaojukeji.know.streaming.km.common.utils.ValidateUtils;
import com.xiaojukeji.know.streaming.km.core.service.broker.BrokerService;
import com.xiaojukeji.know.streaming.km.core.service.health.checkresult.HealthCheckResultService;
import com.xiaojukeji.know.streaming.km.core.service.health.state.HealthStateService;
import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZookeeperService;
import org.apache.commons.collections.CollectionUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import static com.xiaojukeji.know.streaming.km.core.service.version.metrics.BrokerMetricVersionItems.*;
import static com.xiaojukeji.know.streaming.km.core.service.version.metrics.BrokerMetricVersionItems.BROKER_METRIC_HEALTH_STATE;
import static com.xiaojukeji.know.streaming.km.core.service.version.metrics.ClusterMetricVersionItems.*;
import static com.xiaojukeji.know.streaming.km.core.service.version.metrics.GroupMetricVersionItems.*;
import static com.xiaojukeji.know.streaming.km.core.service.version.metrics.GroupMetricVersionItems.GROUP_METRIC_HEALTH_CHECK_TOTAL;
import static com.xiaojukeji.know.streaming.km.core.service.version.metrics.TopicMetricVersionItems.*;
import static com.xiaojukeji.know.streaming.km.core.service.version.metrics.TopicMetricVersionItems.TOPIC_METRIC_HEALTH_CHECK_TOTAL;
import static com.xiaojukeji.know.streaming.km.core.service.version.metrics.ZookeeperMetricVersionItems.*;
@Service
public class HealthStateServiceImpl implements HealthStateService {
@Autowired
private HealthCheckResultService healthCheckResultService;
@Autowired
private ZookeeperService zookeeperService;
@Autowired
private BrokerService brokerService;
@Override
public ClusterMetrics calClusterHealthMetrics(Long clusterPhyId) {
ClusterMetrics metrics = new ClusterMetrics(clusterPhyId);
// 集群维度指标
List<HealthCheckAggResult> resultList = this.getDimensionHealthCheckAggResult(clusterPhyId, HealthCheckDimensionEnum.CLUSTER);
if (ValidateUtils.isEmptyList(resultList)) {
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED_CLUSTER, 0.0f);
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_CLUSTER, 0.0f);
} else {
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED_CLUSTER, this.getHealthCheckPassed(resultList));
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_CLUSTER, (float)resultList.size());
}
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_STATE_CLUSTER, (float)this.calHealthState(resultList).getDimension());
// 获取指标
metrics.putMetric(this.calClusterBrokersHealthMetrics(clusterPhyId).getMetrics());
metrics.putMetric(this.calClusterTopicsHealthMetrics(clusterPhyId).getMetrics());
metrics.putMetric(this.calClusterGroupsHealthMetrics(clusterPhyId).getMetrics());
metrics.putMetric(this.calZookeeperHealthMetrics(clusterPhyId).getMetrics());
// 统计最终结果
Float passed = 0.0f;
passed += metrics.getMetric(ZOOKEEPER_METRIC_HEALTH_CHECK_PASSED);
passed += metrics.getMetric(CLUSTER_METRIC_HEALTH_CHECK_PASSED_TOPICS);
passed += metrics.getMetric(CLUSTER_METRIC_HEALTH_CHECK_PASSED_BROKERS);
passed += metrics.getMetric(CLUSTER_METRIC_HEALTH_CHECK_PASSED_GROUPS);
passed += metrics.getMetric(CLUSTER_METRIC_HEALTH_CHECK_PASSED_CLUSTER);
Float total = 0.0f;
total += metrics.getMetric(ZOOKEEPER_METRIC_HEALTH_CHECK_TOTAL);
total += metrics.getMetric(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_TOPICS);
total += metrics.getMetric(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_BROKERS);
total += metrics.getMetric(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_GROUPS);
total += metrics.getMetric(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_CLUSTER);
// 状态
Float state = 0.0f;
state = Math.max(state, metrics.getMetric(ZOOKEEPER_METRIC_HEALTH_STATE));
state = Math.max(state, metrics.getMetric(CLUSTER_METRIC_HEALTH_STATE_TOPICS));
state = Math.max(state, metrics.getMetric(CLUSTER_METRIC_HEALTH_STATE_BROKERS));
state = Math.max(state, metrics.getMetric(CLUSTER_METRIC_HEALTH_STATE_GROUPS));
state = Math.max(state, metrics.getMetric(CLUSTER_METRIC_HEALTH_STATE_CLUSTER));
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED, passed);
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL, total);
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_STATE, state);
return metrics;
}
@Override
public BrokerMetrics calBrokerHealthMetrics(Long clusterPhyId, Integer brokerId) {
List<HealthScoreResult> healthScoreResultList = this.getResHealthResult(clusterPhyId, HealthCheckDimensionEnum.BROKER.getDimension(), String.valueOf(brokerId));
BrokerMetrics metrics = new BrokerMetrics(clusterPhyId, brokerId);
if (ValidateUtils.isEmptyList(healthScoreResultList)) {
metrics.getMetrics().put(BROKER_METRIC_HEALTH_STATE, (float)HealthStateEnum.GOOD.getDimension());
metrics.getMetrics().put(BROKER_METRIC_HEALTH_CHECK_PASSED, 0.0f);
metrics.getMetrics().put(BROKER_METRIC_HEALTH_CHECK_TOTAL, 0.0f);
} else {
metrics.getMetrics().put(BROKER_METRIC_HEALTH_CHECK_PASSED, getHealthCheckResultPassed(healthScoreResultList));
metrics.getMetrics().put(BROKER_METRIC_HEALTH_CHECK_TOTAL, Float.valueOf(healthScoreResultList.size()));
// 计算健康状态
Broker broker = brokerService.getBrokerFromCacheFirst(clusterPhyId, brokerId);
if (broker == null) {
// DB中不存在则默认是存活的
metrics.getMetrics().put(BROKER_METRIC_HEALTH_STATE, (float)HealthStateEnum.GOOD.getDimension());
} else if (!broker.alive()) {
metrics.getMetrics().put(BROKER_METRIC_HEALTH_STATE, (float)HealthStateEnum.DEAD.getDimension());
} else {
metrics.getMetrics().put(BROKER_METRIC_HEALTH_STATE, (float)this.calHealthScoreResultState(healthScoreResultList).getDimension());
}
}
return metrics;
}
@Override
public TopicMetrics calTopicHealthMetrics(Long clusterPhyId, String topicName) {
List<HealthScoreResult> healthScoreResultList = this.getResHealthResult(clusterPhyId, HealthCheckDimensionEnum.TOPIC.getDimension(), topicName);
TopicMetrics metrics = new TopicMetrics(topicName, clusterPhyId,true);
if (ValidateUtils.isEmptyList(healthScoreResultList)) {
metrics.getMetrics().put(TOPIC_METRIC_HEALTH_STATE, (float)HealthStateEnum.GOOD.getDimension());
metrics.getMetrics().put(TOPIC_METRIC_HEALTH_CHECK_PASSED, 0.0f);
metrics.getMetrics().put(TOPIC_METRIC_HEALTH_CHECK_TOTAL, 0.0f);
} else {
metrics.getMetrics().put(TOPIC_METRIC_HEALTH_STATE, (float)this.calHealthScoreResultState(healthScoreResultList).getDimension());
metrics.getMetrics().put(TOPIC_METRIC_HEALTH_CHECK_PASSED, this.getHealthCheckResultPassed(healthScoreResultList));
metrics.getMetrics().put(TOPIC_METRIC_HEALTH_CHECK_TOTAL, Float.valueOf(healthScoreResultList.size()));
}
return metrics;
}
@Override
public GroupMetrics calGroupHealthMetrics(Long clusterPhyId, String groupName) {
List<HealthScoreResult> healthScoreResultList = this.getResHealthResult(clusterPhyId, HealthCheckDimensionEnum.GROUP.getDimension(), groupName);
GroupMetrics metrics = new GroupMetrics(clusterPhyId, groupName, true);
if (ValidateUtils.isEmptyList(healthScoreResultList)) {
metrics.getMetrics().put(GROUP_METRIC_HEALTH_STATE, (float)HealthStateEnum.GOOD.getDimension());
metrics.getMetrics().put(GROUP_METRIC_HEALTH_CHECK_PASSED, 0.0f);
metrics.getMetrics().put(GROUP_METRIC_HEALTH_CHECK_TOTAL, 0.0f);
} else {
metrics.getMetrics().put(GROUP_METRIC_HEALTH_STATE, (float)this.calHealthScoreResultState(healthScoreResultList).getDimension());
metrics.getMetrics().put(GROUP_METRIC_HEALTH_CHECK_PASSED, getHealthCheckResultPassed(healthScoreResultList));
metrics.getMetrics().put(GROUP_METRIC_HEALTH_CHECK_TOTAL, Float.valueOf(healthScoreResultList.size()));
}
return metrics;
}
@Override
public ZookeeperMetrics calZookeeperHealthMetrics(Long clusterPhyId) {
List<HealthCheckAggResult> resultList = this.getDimensionHealthCheckAggResult(clusterPhyId, HealthCheckDimensionEnum.ZOOKEEPER);
ZookeeperMetrics metrics = new ZookeeperMetrics(clusterPhyId);
if (ValidateUtils.isEmptyList(resultList)) {
metrics.getMetrics().put(ZOOKEEPER_METRIC_HEALTH_CHECK_PASSED, 0.0f);
metrics.getMetrics().put(ZOOKEEPER_METRIC_HEALTH_CHECK_TOTAL, 0.0f);
} else {
metrics.getMetrics().put(ZOOKEEPER_METRIC_HEALTH_CHECK_PASSED, this.getHealthCheckPassed(resultList));
metrics.getMetrics().put(ZOOKEEPER_METRIC_HEALTH_CHECK_TOTAL, (float)resultList.size());
}
if (zookeeperService.allServerDown(clusterPhyId)) {
// 所有服务挂掉
metrics.getMetrics().put(ZOOKEEPER_METRIC_HEALTH_STATE, (float)HealthStateEnum.DEAD.getDimension());
return metrics;
}
if (zookeeperService.existServerDown(clusterPhyId)) {
// 存在服务挂掉
metrics.getMetrics().put(ZOOKEEPER_METRIC_HEALTH_STATE, (float)HealthStateEnum.POOR.getDimension());
return metrics;
}
// 服务未挂时,依据检查结果计算状态
metrics.getMetrics().put(ZOOKEEPER_METRIC_HEALTH_STATE, (float)this.calHealthState(resultList).getDimension());
return metrics;
}
@Override
public List<HealthScoreResult> getClusterHealthResult(Long clusterPhyId) {
List<HealthCheckResultPO> poList = healthCheckResultService.getClusterHealthCheckResult(clusterPhyId);
// <检查项,<检查结果>>
Map<String, List<HealthCheckResultPO>> checkResultMap = new HashMap<>();
for (HealthCheckResultPO po: poList) {
checkResultMap.putIfAbsent(po.getConfigName(), new ArrayList<>());
checkResultMap.get(po.getConfigName()).add(po);
}
Map<String, BaseClusterHealthConfig> configMap = healthCheckResultService.getClusterHealthConfig(clusterPhyId);
List<HealthScoreResult> healthScoreResultList = new ArrayList<>();
for (HealthCheckNameEnum nameEnum: HealthCheckNameEnum.values()) {
BaseClusterHealthConfig baseConfig = configMap.get(nameEnum.getConfigName());
if (baseConfig == null) {
continue;
}
healthScoreResultList.add(new HealthScoreResult(
nameEnum,
baseConfig,
checkResultMap.getOrDefault(nameEnum.getConfigName(), new ArrayList<>()))
);
}
return healthScoreResultList;
}
@Override
public List<HealthScoreResult> getDimensionHealthResult(Long clusterPhyId, HealthCheckDimensionEnum dimensionEnum) {
List<HealthCheckResultPO> poList = healthCheckResultService.getClusterResourcesHealthCheckResult(clusterPhyId, dimensionEnum.getDimension());
// <检查项,<通过的数量,不通过的数量>>
Map<String, List<HealthCheckResultPO>> checkResultMap = new HashMap<>();
for (HealthCheckResultPO po: poList) {
checkResultMap.putIfAbsent(po.getConfigName(), new ArrayList<>());
checkResultMap.get(po.getConfigName()).add(po);
}
Map<String, BaseClusterHealthConfig> configMap = healthCheckResultService.getClusterHealthConfig(clusterPhyId);
List<HealthScoreResult> healthScoreResultList = new ArrayList<>();
for (HealthCheckNameEnum nameEnum: HealthCheckNameEnum.getByDimension(dimensionEnum)) {
BaseClusterHealthConfig baseConfig = configMap.get(nameEnum.getConfigName());
if (baseConfig == null) {
continue;
}
healthScoreResultList.add(new HealthScoreResult(nameEnum, baseConfig, checkResultMap.getOrDefault(nameEnum.getConfigName(), new ArrayList<>())));
}
return healthScoreResultList;
}
@Override
public List<HealthScoreResult> getResHealthResult(Long clusterPhyId, Integer dimension, String resNme) {
List<HealthCheckResultPO> poList = healthCheckResultService.getResHealthCheckResult(clusterPhyId, dimension, resNme);
Map<String, List<HealthCheckResultPO>> checkResultMap = new HashMap<>();
for (HealthCheckResultPO po: poList) {
checkResultMap.putIfAbsent(po.getConfigName(), new ArrayList<>());
checkResultMap.get(po.getConfigName()).add(po);
}
Map<String, BaseClusterHealthConfig> configMap = healthCheckResultService.getClusterHealthConfig(clusterPhyId);
List<HealthScoreResult> healthScoreResultList = new ArrayList<>();
for (HealthCheckNameEnum nameEnum: HealthCheckNameEnum.getByDimensionCode(dimension)) {
BaseClusterHealthConfig baseConfig = configMap.get(nameEnum.getConfigName());
if (baseConfig == null) {
continue;
}
healthScoreResultList.add(new HealthScoreResult(nameEnum, baseConfig, checkResultMap.getOrDefault(nameEnum.getConfigName(), new ArrayList<>())));
}
return healthScoreResultList;
}
/**************************************************** private method ****************************************************/
private ClusterMetrics calClusterTopicsHealthMetrics(Long clusterPhyId) {
List<HealthCheckAggResult> resultList = this.getDimensionHealthCheckAggResult(clusterPhyId, HealthCheckDimensionEnum.TOPIC);
ClusterMetrics metrics = new ClusterMetrics(clusterPhyId);
if (ValidateUtils.isEmptyList(resultList)) {
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED_TOPICS, 0.0f);
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_TOPICS, 0.0f);
} else {
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED_TOPICS, this.getHealthCheckPassed(resultList));
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_TOPICS, (float)resultList.size());
}
// 服务未挂时,依据检查结果计算状态
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_STATE_TOPICS, (float)this.calHealthState(resultList).getDimension());
return metrics;
}
private ClusterMetrics calClusterGroupsHealthMetrics(Long clusterPhyId) {
List<HealthCheckAggResult> resultList = this.getDimensionHealthCheckAggResult(clusterPhyId, HealthCheckDimensionEnum.GROUP);
ClusterMetrics metrics = new ClusterMetrics(clusterPhyId);
if (ValidateUtils.isEmptyList(resultList)) {
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED_GROUPS, 0.0f);
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_GROUPS, 0.0f);
} else {
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED_GROUPS, this.getHealthCheckPassed(resultList));
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_GROUPS, (float)resultList.size());
}
// 服务未挂时,依据检查结果计算状态
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_STATE_GROUPS, (float)this.calHealthState(resultList).getDimension());
return metrics;
}
private ClusterMetrics calClusterBrokersHealthMetrics(Long clusterPhyId) {
List<HealthCheckAggResult> resultList = this.getDimensionHealthCheckAggResult(clusterPhyId, HealthCheckDimensionEnum.BROKER);
ClusterMetrics metrics = new ClusterMetrics(clusterPhyId);
if (ValidateUtils.isEmptyList(resultList)) {
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED_BROKERS, 0.0f);
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_BROKERS, 0.0f);
} else {
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_PASSED_BROKERS, this.getHealthCheckPassed(resultList));
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_BROKERS, (float)resultList.size());
}
if (brokerService.allServerDown(clusterPhyId)) {
// 所有服务挂掉
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_STATE_BROKERS, (float)HealthStateEnum.DEAD.getDimension());
return metrics;
}
if (brokerService.existServerDown(clusterPhyId)) {
// 存在服务挂掉
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_STATE_BROKERS, (float)HealthStateEnum.POOR.getDimension());
return metrics;
}
// 服务未挂时,依据检查结果计算状态
metrics.getMetrics().put(CLUSTER_METRIC_HEALTH_STATE_BROKERS, (float)this.calHealthState(resultList).getDimension());
return metrics;
}
private List<HealthCheckAggResult> getDimensionHealthCheckAggResult(Long clusterPhyId, HealthCheckDimensionEnum dimensionEnum) {
List<HealthCheckResultPO> poList = healthCheckResultService.getClusterResourcesHealthCheckResult(clusterPhyId, dimensionEnum.getDimension());
Map<String /*检查名*/, List<HealthCheckResultPO> /*检查结果列表*/> groupByCheckNamePOMap = new HashMap<>();
for (HealthCheckResultPO po: poList) {
groupByCheckNamePOMap.putIfAbsent(po.getConfigName(), new ArrayList<>());
groupByCheckNamePOMap.get(po.getConfigName()).add(po);
}
List<HealthCheckAggResult> stateList = new ArrayList<>();
for (HealthCheckNameEnum nameEnum: HealthCheckNameEnum.getByDimension(dimensionEnum)) {
stateList.add(new HealthCheckAggResult(nameEnum, groupByCheckNamePOMap.getOrDefault(nameEnum.getConfigName(), new ArrayList<>())));
}
return stateList;
}
private float getHealthCheckPassed(List<HealthCheckAggResult> resultList){
if(ValidateUtils.isEmptyList(resultList)) {
return 0f;
}
return Float.valueOf(resultList.stream().filter(elem -> elem.getPassed()).count());
}
private HealthStateEnum calHealthState(List<HealthCheckAggResult> resultList) {
if(ValidateUtils.isEmptyList(resultList)) {
return HealthStateEnum.GOOD;
}
boolean existNotPassed = false;
for (HealthCheckAggResult aggResult: resultList) {
if (aggResult.getCheckNameEnum().isAvailableChecker() && !aggResult.getPassed()) {
return HealthStateEnum.POOR;
}
if (!aggResult.getPassed()) {
existNotPassed = true;
}
}
return existNotPassed? HealthStateEnum.MEDIUM: HealthStateEnum.GOOD;
}
private float getHealthCheckResultPassed(List<HealthScoreResult> healthScoreResultList){
if(CollectionUtils.isEmpty(healthScoreResultList)){return 0f;}
return Float.valueOf(healthScoreResultList.stream().filter(elem -> elem.getPassed()).count());
}
private HealthStateEnum calHealthScoreResultState(List<HealthScoreResult> resultList) {
if(ValidateUtils.isEmptyList(resultList)) {
return HealthStateEnum.GOOD;
}
boolean existNotPassed = false;
for (HealthScoreResult aggResult: resultList) {
if (aggResult.getCheckNameEnum().isAvailableChecker() && !aggResult.getPassed()) {
return HealthStateEnum.POOR;
}
if (!aggResult.getPassed()) {
existNotPassed = true;
}
}
return existNotPassed? HealthStateEnum.MEDIUM: HealthStateEnum.GOOD;
}
}

View File

@@ -31,7 +31,7 @@ import com.xiaojukeji.know.streaming.km.common.utils.ConvertUtil;
import com.xiaojukeji.know.streaming.km.common.utils.ValidateUtils; import com.xiaojukeji.know.streaming.km.common.utils.ValidateUtils;
import com.xiaojukeji.know.streaming.km.core.cache.CollectedMetricsLocalCache; import com.xiaojukeji.know.streaming.km.core.cache.CollectedMetricsLocalCache;
import com.xiaojukeji.know.streaming.km.core.service.broker.BrokerService; import com.xiaojukeji.know.streaming.km.core.service.broker.BrokerService;
import com.xiaojukeji.know.streaming.km.core.service.health.score.HealthScoreService; import com.xiaojukeji.know.streaming.km.core.service.health.state.HealthStateService;
import com.xiaojukeji.know.streaming.km.core.service.partition.PartitionMetricService; import com.xiaojukeji.know.streaming.km.core.service.partition.PartitionMetricService;
import com.xiaojukeji.know.streaming.km.core.service.topic.TopicMetricService; import com.xiaojukeji.know.streaming.km.core.service.topic.TopicMetricService;
import com.xiaojukeji.know.streaming.km.core.service.topic.TopicService; import com.xiaojukeji.know.streaming.km.core.service.topic.TopicService;
@@ -69,7 +69,7 @@ public class TopicMetricServiceImpl extends BaseMetricService implements TopicMe
public static final String TOPIC_METHOD_GET_REPLICAS_COUNT = "getReplicasCount"; public static final String TOPIC_METHOD_GET_REPLICAS_COUNT = "getReplicasCount";
@Autowired @Autowired
private HealthScoreService healthScoreService; private HealthStateService healthStateService;
@Autowired @Autowired
private KafkaJMXClient kafkaJMXClient; private KafkaJMXClient kafkaJMXClient;
@@ -394,7 +394,7 @@ public class TopicMetricServiceImpl extends BaseMetricService implements TopicMe
String topic = topicMetricParam.getTopic(); String topic = topicMetricParam.getTopic();
Long clusterId = topicMetricParam.getClusterId(); Long clusterId = topicMetricParam.getClusterId();
TopicMetrics topicMetric = healthScoreService.calTopicHealthScore(clusterId, topic); TopicMetrics topicMetric = healthStateService.calTopicHealthMetrics(clusterId, topic);
return Result.buildSuc(Arrays.asList(topicMetric)); return Result.buildSuc(Arrays.asList(topicMetric));
} }

View File

@@ -16,7 +16,7 @@ import static com.xiaojukeji.know.streaming.km.core.service.broker.impl.BrokerMe
@Component @Component
public class BrokerMetricVersionItems extends BaseMetricVersionMetric { public class BrokerMetricVersionItems extends BaseMetricVersionMetric {
public static final String BROKER_METRIC_HEALTH_SCORE = "HealthScore"; public static final String BROKER_METRIC_HEALTH_STATE = "HealthState";
public static final String BROKER_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed"; public static final String BROKER_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed";
public static final String BROKER_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal"; public static final String BROKER_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal";
public static final String BROKER_METRIC_TOTAL_REQ_QUEUE = "TotalRequestQueueSize"; public static final String BROKER_METRIC_TOTAL_REQ_QUEUE = "TotalRequestQueueSize";
@@ -57,8 +57,8 @@ public class BrokerMetricVersionItems extends BaseMetricVersionMetric {
// HealthScore 指标 // HealthScore 指标
items.add(buildAllVersionsItem() items.add(buildAllVersionsItem()
.name(BROKER_METRIC_HEALTH_SCORE).unit("").desc("健康").category(CATEGORY_HEALTH) .name(BROKER_METRIC_HEALTH_STATE).unit("0:好 1:中 2:差 3:宕机").desc("健康状态(0:好 1:中 2:差 3:宕机)").category(CATEGORY_HEALTH)
.extendMethod(BROKER_METHOD_GET_HEALTH_SCORE)); .extendMethod( BROKER_METHOD_GET_HEALTH_SCORE ));
items.add(buildAllVersionsItem() items.add(buildAllVersionsItem()
.name(BROKER_METRIC_HEALTH_CHECK_PASSED ).unit("").desc("健康检查通过数").category(CATEGORY_HEALTH) .name(BROKER_METRIC_HEALTH_CHECK_PASSED ).unit("").desc("健康检查通过数").category(CATEGORY_HEALTH)

View File

@@ -19,36 +19,41 @@ import static com.xiaojukeji.know.streaming.km.core.service.cluster.impl.Cluster
*/ */
@Component @Component
public class ClusterMetricVersionItems extends BaseMetricVersionMetric { public class ClusterMetricVersionItems extends BaseMetricVersionMetric {
/** /**
* 健康分 * 整体的健康指标
*/
public static final String CLUSTER_METRIC_HEALTH_SCORE = "HealthScore";
public static final String CLUSTER_METRIC_HEALTH_SCORE_TOPICS = "HealthScore_Topics";
public static final String CLUSTER_METRIC_HEALTH_SCORE_BROKERS = "HealthScore_Brokers";
public static final String CLUSTER_METRIC_HEALTH_SCORE_GROUPS = "HealthScore_Groups";
public static final String CLUSTER_METRIC_HEALTH_SCORE_CLUSTER = "HealthScore_Cluster";
/**
* 健康巡检
*/ */
public static final String CLUSTER_METRIC_HEALTH_STATE = "HealthState";
public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed"; public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed";
public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal"; public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal";
/**
* Topics健康指标
*/
public static final String CLUSTER_METRIC_HEALTH_STATE_TOPICS = "HealthState_Topics";
public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED_TOPICS = "HealthCheckPassed_Topics"; public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED_TOPICS = "HealthCheckPassed_Topics";
public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL_TOPICS = "HealthCheckTotal_Topics"; public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL_TOPICS = "HealthCheckTotal_Topics";
/**
* Brokers健康指标
*/
public static final String CLUSTER_METRIC_HEALTH_STATE_BROKERS = "HealthState_Brokers";
public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED_BROKERS = "HealthCheckPassed_Brokers"; public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED_BROKERS = "HealthCheckPassed_Brokers";
public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL_BROKERS = "HealthCheckTotal_Brokers"; public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL_BROKERS = "HealthCheckTotal_Brokers";
/**
* Groups健康指标
*/
public static final String CLUSTER_METRIC_HEALTH_STATE_GROUPS = "HealthState_Groups";
public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED_GROUPS = "HealthCheckPassed_Groups"; public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED_GROUPS = "HealthCheckPassed_Groups";
public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL_GROUPS = "HealthCheckTotal_Groups"; public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL_GROUPS = "HealthCheckTotal_Groups";
/**
* Cluster健康指标
*/
public static final String CLUSTER_METRIC_HEALTH_STATE_CLUSTER = "HealthState_Cluster";
public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED_CLUSTER = "HealthCheckPassed_Cluster"; public static final String CLUSTER_METRIC_HEALTH_CHECK_PASSED_CLUSTER = "HealthCheckPassed_Cluster";
public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL_CLUSTER = "HealthCheckTotal_Cluster"; public static final String CLUSTER_METRIC_HEALTH_CHECK_TOTAL_CLUSTER = "HealthCheckTotal_Cluster";
public static final String CLUSTER_METRIC_TOTAL_REQ_QUEUE_SIZE = "TotalRequestQueueSize"; public static final String CLUSTER_METRIC_TOTAL_REQ_QUEUE_SIZE = "TotalRequestQueueSize";
public static final String CLUSTER_METRIC_TOTAL_RES_QUEUE_SIZE = "TotalResponseQueueSize"; public static final String CLUSTER_METRIC_TOTAL_RES_QUEUE_SIZE = "TotalResponseQueueSize";
public static final String CLUSTER_METRIC_EVENT_QUEUE_SIZE = "EventQueueSize"; public static final String CLUSTER_METRIC_EVENT_QUEUE_SIZE = "EventQueueSize";
@@ -113,64 +118,64 @@ public class ClusterMetricVersionItems extends BaseMetricVersionMetric {
// HealthScore 指标 // HealthScore 指标
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_SCORE).unit("").desc("集群总体的健康分").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_STATE).unit("0:好 1:中 2:差 3:宕").desc("集群健康状态(0:好 1:中 2:差 3:宕)").category(CATEGORY_HEALTH)
.extendMethod(CLUSTER_METHOD_GET_HEALTH_SCORE)); .extendMethod(CLUSTER_METHOD_GET_HEALTH_METRICS));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_PASSED).unit("").desc("集群总体健康检查通过数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_PASSED).unit("").desc("集群总体健康检查通过数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL).unit("").desc("集群总体健康检查总数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL).unit("").desc("集群总体健康检查总数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_SCORE_TOPICS).unit("").desc("集群Topics健康").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_STATE_TOPICS).unit("0:好 1:中 2:差 3:宕机").desc("集群Topics健康状态").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod(CLUSTER_METHOD_GET_HEALTH_METRICS));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_PASSED_TOPICS).unit("").desc("集群Topics健康检查通过数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_PASSED_TOPICS).unit("").desc("集群Topics健康检查通过数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_TOPICS).unit("").desc("集群Topics健康检查总数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_TOPICS).unit("").desc("集群Topics健康检查总数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_SCORE_BROKERS).unit("").desc("集群Brokers健康").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_STATE_BROKERS).unit("0:好 1:中 2:差 3:宕机").desc("集群Brokers健康状态").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod(CLUSTER_METHOD_GET_HEALTH_METRICS));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_PASSED_BROKERS).unit("").desc("集群Brokers健康检查通过数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_PASSED_BROKERS).unit("").desc("集群Brokers健康检查通过数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_BROKERS).unit("").desc("集群Brokers健康检查总数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_BROKERS).unit("").desc("集群Brokers健康检查总数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_SCORE_GROUPS).unit("").desc("集群Groups健康").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_STATE_GROUPS).unit("0:好 1:中 2:差 3:宕机").desc("集群Groups健康状态").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod(CLUSTER_METHOD_GET_HEALTH_METRICS));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_PASSED_GROUPS).unit("").desc("集群Groups健康检查通过数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_PASSED_GROUPS).unit("").desc("集群Groups健康检查通过数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_GROUPS).unit("").desc("集群Groups健康检查总数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_GROUPS).unit("").desc("集群Groups健康检查总数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_SCORE_CLUSTER).unit("").desc("集群自身健康").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_STATE_CLUSTER).unit("0:好 1:中 2:差 3:宕机").desc("集群自身健康状态").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod(CLUSTER_METHOD_GET_HEALTH_METRICS));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_PASSED_CLUSTER).unit("").desc("集群自身健康检查通过数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_PASSED_CLUSTER).unit("").desc("集群自身健康检查通过数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_CLUSTER).unit("").desc("集群自身健康检查总数").category(CATEGORY_HEALTH) .name(CLUSTER_METRIC_HEALTH_CHECK_TOTAL_CLUSTER).unit("").desc("集群自身健康检查总数").category(CATEGORY_HEALTH)
.extendMethod( CLUSTER_METHOD_GET_HEALTH_SCORE )); .extendMethod( CLUSTER_METHOD_GET_HEALTH_METRICS ));
// TotalRequestQueueSize 指标 // TotalRequestQueueSize 指标
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()

View File

@@ -7,13 +7,14 @@ import org.springframework.stereotype.Component;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import static com.xiaojukeji.know.streaming.km.common.bean.entity.version.VersionMetricControlItem.CATEGORY_HEALTH;
import static com.xiaojukeji.know.streaming.km.common.enums.version.VersionItemTypeEnum.METRIC_GROUP; import static com.xiaojukeji.know.streaming.km.common.enums.version.VersionItemTypeEnum.METRIC_GROUP;
import static com.xiaojukeji.know.streaming.km.core.service.group.impl.GroupMetricServiceImpl.*; import static com.xiaojukeji.know.streaming.km.core.service.group.impl.GroupMetricServiceImpl.*;
@Component @Component
public class GroupMetricVersionItems extends BaseMetricVersionMetric { public class GroupMetricVersionItems extends BaseMetricVersionMetric {
public static final String GROUP_METRIC_HEALTH_SCORE = "HealthScore"; public static final String GROUP_METRIC_HEALTH_STATE = "HealthState";
public static final String GROUP_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed"; public static final String GROUP_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed";
public static final String GROUP_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal"; public static final String GROUP_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal";
public static final String GROUP_METRIC_OFFSET_CONSUMED = "OffsetConsumed"; public static final String GROUP_METRIC_OFFSET_CONSUMED = "OffsetConsumed";
@@ -33,15 +34,15 @@ public class GroupMetricVersionItems extends BaseMetricVersionMetric {
// HealthScore 指标 // HealthScore 指标
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(GROUP_METRIC_HEALTH_SCORE).unit("").desc("健康") .name(GROUP_METRIC_HEALTH_STATE).unit("0:好 1:中 2:差 3:宕机").desc("健康状态(0:好 1:中 2:差 3:宕机)").category(CATEGORY_HEALTH)
.extendMethod( GROUP_METHOD_GET_HEALTH_SCORE )); .extendMethod( GROUP_METHOD_GET_HEALTH_SCORE ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(GROUP_METRIC_HEALTH_CHECK_PASSED ).unit("").desc("健康检查通过数") .name(GROUP_METRIC_HEALTH_CHECK_PASSED ).unit("").desc("健康检查通过数").category(CATEGORY_HEALTH)
.extendMethod( GROUP_METHOD_GET_HEALTH_SCORE )); .extendMethod( GROUP_METHOD_GET_HEALTH_SCORE ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(GROUP_METRIC_HEALTH_CHECK_TOTAL ).unit("").desc("健康检查总数") .name(GROUP_METRIC_HEALTH_CHECK_TOTAL ).unit("").desc("健康检查总数").category(CATEGORY_HEALTH)
.extendMethod( GROUP_METHOD_GET_HEALTH_SCORE )); .extendMethod( GROUP_METHOD_GET_HEALTH_SCORE ));
// OffsetConsumed 指标 // OffsetConsumed 指标

View File

@@ -16,9 +16,10 @@ import static com.xiaojukeji.know.streaming.km.core.service.topic.impl.TopicMetr
@Component @Component
public class TopicMetricVersionItems extends BaseMetricVersionMetric { public class TopicMetricVersionItems extends BaseMetricVersionMetric {
public static final String TOPIC_METRIC_HEALTH_SCORE = "HealthScore"; public static final String TOPIC_METRIC_HEALTH_STATE = "HealthState";
public static final String TOPIC_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed"; public static final String TOPIC_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed";
public static final String TOPIC_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal"; public static final String TOPIC_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal";
public static final String TOPIC_METRIC_TOTAL_PRODUCE_REQUESTS = "TotalProduceRequests"; public static final String TOPIC_METRIC_TOTAL_PRODUCE_REQUESTS = "TotalProduceRequests";
public static final String TOPIC_METRIC_BYTES_REJECTED = "BytesRejected"; public static final String TOPIC_METRIC_BYTES_REJECTED = "BytesRejected";
public static final String TOPIC_METRIC_FAILED_FETCH_REQ = "FailedFetchRequests"; public static final String TOPIC_METRIC_FAILED_FETCH_REQ = "FailedFetchRequests";
@@ -47,7 +48,7 @@ public class TopicMetricVersionItems extends BaseMetricVersionMetric {
// HealthScore 指标 // HealthScore 指标
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()
.name(TOPIC_METRIC_HEALTH_SCORE).unit("").desc("健康").category(CATEGORY_HEALTH) .name(TOPIC_METRIC_HEALTH_STATE).unit("0:好 1:中 2:差 3:宕机").desc("健康状态(0:好 1:中 2:差 3:宕机)").category(CATEGORY_HEALTH)
.extendMethod( TOPIC_METHOD_GET_HEALTH_SCORE )); .extendMethod( TOPIC_METHOD_GET_HEALTH_SCORE ));
itemList.add(buildAllVersionsItem() itemList.add(buildAllVersionsItem()

View File

@@ -15,6 +15,12 @@ import static com.xiaojukeji.know.streaming.km.core.service.zookeeper.impl.Zooke
@Component @Component
public class ZookeeperMetricVersionItems extends BaseMetricVersionMetric { public class ZookeeperMetricVersionItems extends BaseMetricVersionMetric {
/**
* 健康状态
*/
public static final String ZOOKEEPER_METRIC_HEALTH_STATE = "HealthState";
public static final String ZOOKEEPER_METRIC_HEALTH_CHECK_PASSED = "HealthCheckPassed";
public static final String ZOOKEEPER_METRIC_HEALTH_CHECK_TOTAL = "HealthCheckTotal";
/** /**
* 性能 * 性能
@@ -23,7 +29,7 @@ public class ZookeeperMetricVersionItems extends BaseMetricVersionMetric {
public static final String ZOOKEEPER_METRIC_MIN_REQUEST_LATENCY = "MinRequestLatency"; public static final String ZOOKEEPER_METRIC_MIN_REQUEST_LATENCY = "MinRequestLatency";
public static final String ZOOKEEPER_METRIC_MAX_REQUEST_LATENCY = "MaxRequestLatency"; public static final String ZOOKEEPER_METRIC_MAX_REQUEST_LATENCY = "MaxRequestLatency";
public static final String ZOOKEEPER_METRIC_OUTSTANDING_REQUESTS = "OutstandingRequests"; public static final String ZOOKEEPER_METRIC_OUTSTANDING_REQUESTS = "OutstandingRequests";
public static final String ZOOKEEPER_METRIC_NODE_COUNT = "NodeCount"; public static final String ZOOKEEPER_METRIC_NODE_COUNT = "ZnodeCount";
public static final String ZOOKEEPER_METRIC_WATCH_COUNT = "WatchCount"; public static final String ZOOKEEPER_METRIC_WATCH_COUNT = "WatchCount";
public static final String ZOOKEEPER_METRIC_NUM_ALIVE_CONNECTIONS = "NumAliveConnections"; public static final String ZOOKEEPER_METRIC_NUM_ALIVE_CONNECTIONS = "NumAliveConnections";
public static final String ZOOKEEPER_METRIC_PACKETS_RECEIVED = "PacketsReceived"; public static final String ZOOKEEPER_METRIC_PACKETS_RECEIVED = "PacketsReceived";
@@ -52,6 +58,20 @@ public class ZookeeperMetricVersionItems extends BaseMetricVersionMetric {
public List<VersionMetricControlItem> init(){ public List<VersionMetricControlItem> init(){
List<VersionMetricControlItem> items = new ArrayList<>(); List<VersionMetricControlItem> items = new ArrayList<>();
// 健康状态
items.add(buildAllVersionsItem()
.name(ZOOKEEPER_METRIC_HEALTH_STATE).unit("0:好 1:中 2:差 3:宕机").desc("健康状态(0:好 1:中 2:差 3:宕机)").category(CATEGORY_HEALTH)
.extendMethod(ZOOKEEPER_METHOD_GET_METRIC_FROM_HEALTH_SERVICE));
items.add(buildAllVersionsItem()
.name(ZOOKEEPER_METRIC_HEALTH_CHECK_PASSED).unit("").desc("健康巡检通过数").category(CATEGORY_HEALTH)
.extendMethod(ZOOKEEPER_METHOD_GET_METRIC_FROM_HEALTH_SERVICE));
items.add(buildAllVersionsItem()
.name(ZOOKEEPER_METRIC_HEALTH_CHECK_TOTAL).unit("").desc("健康巡检总数").category(CATEGORY_HEALTH)
.extendMethod(ZOOKEEPER_METHOD_GET_METRIC_FROM_HEALTH_SERVICE));
// 性能指标 // 性能指标
items.add(buildAllVersionsItem() items.add(buildAllVersionsItem()
.name(ZOOKEEPER_METRIC_AVG_REQUEST_LATENCY).unit("ms").desc("平均响应延迟").category(CATEGORY_PERFORMANCE) .name(ZOOKEEPER_METRIC_AVG_REQUEST_LATENCY).unit("ms").desc("平均响应延迟").category(CATEGORY_PERFORMANCE)

View File

@@ -29,6 +29,7 @@ import com.xiaojukeji.know.streaming.km.common.bean.po.metrice.ZookeeperMetricPO
import com.xiaojukeji.know.streaming.km.common.utils.zookeeper.FourLetterWordUtil; import com.xiaojukeji.know.streaming.km.common.utils.zookeeper.FourLetterWordUtil;
import com.xiaojukeji.know.streaming.km.core.cache.ZookeeperLocalCache; import com.xiaojukeji.know.streaming.km.core.cache.ZookeeperLocalCache;
import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterPhyService; import com.xiaojukeji.know.streaming.km.core.service.cluster.ClusterPhyService;
import com.xiaojukeji.know.streaming.km.core.service.health.state.HealthStateService;
import com.xiaojukeji.know.streaming.km.core.service.version.BaseMetricService; import com.xiaojukeji.know.streaming.km.core.service.version.BaseMetricService;
import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZookeeperMetricService; import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZookeeperMetricService;
import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZookeeperService; import com.xiaojukeji.know.streaming.km.core.service.zookeeper.ZookeeperService;
@@ -68,6 +69,9 @@ public class ZookeeperMetricServiceImpl extends BaseMetricService implements Zoo
@Autowired @Autowired
private KafkaJMXClient kafkaJMXClient; private KafkaJMXClient kafkaJMXClient;
@Autowired
private HealthStateService healthStateService;
@Override @Override
protected VersionItemTypeEnum getVersionItemType() { protected VersionItemTypeEnum getVersionItemType() {
return VersionItemTypeEnum.METRIC_ZOOKEEPER; return VersionItemTypeEnum.METRIC_ZOOKEEPER;
@@ -84,6 +88,7 @@ public class ZookeeperMetricServiceImpl extends BaseMetricService implements Zoo
registerVCHandler( ZOOKEEPER_METHOD_GET_METRIC_FROM_MONITOR_CMD, this::getMetricFromMonitorCmd); registerVCHandler( ZOOKEEPER_METHOD_GET_METRIC_FROM_MONITOR_CMD, this::getMetricFromMonitorCmd);
registerVCHandler( ZOOKEEPER_METHOD_GET_METRIC_FROM_SERVER_CMD, this::getMetricFromServerCmd); registerVCHandler( ZOOKEEPER_METHOD_GET_METRIC_FROM_SERVER_CMD, this::getMetricFromServerCmd);
registerVCHandler( ZOOKEEPER_METHOD_GET_METRIC_FROM_KAFKA_BY_JMX, this::getMetricFromKafkaByJMX); registerVCHandler( ZOOKEEPER_METHOD_GET_METRIC_FROM_KAFKA_BY_JMX, this::getMetricFromKafkaByJMX);
registerVCHandler( ZOOKEEPER_METHOD_GET_METRIC_FROM_HEALTH_SERVICE, this::getMetricFromHealthService);
} }
@Override @Override
@@ -302,4 +307,13 @@ public class ZookeeperMetricServiceImpl extends BaseMetricService implements Zoo
return Result.buildFailure(VC_JMX_CONNECT_ERROR); return Result.buildFailure(VC_JMX_CONNECT_ERROR);
} }
} }
private Result<ZookeeperMetrics> getMetricFromHealthService(VersionItemParam metricParam) {
ZookeeperMetricParam param = (ZookeeperMetricParam)metricParam;
String metricName = param.getMetricName();
Long clusterPhyId = param.getClusterPhyId();
return Result.buildSuc(healthStateService.calZookeeperHealthMetrics(clusterPhyId));
}
} }

View File

@@ -10,7 +10,7 @@ import com.xiaojukeji.know.streaming.km.common.converter.HealthScoreVOConverter;
import com.xiaojukeji.know.streaming.km.common.enums.config.ConfigGroupEnum; import com.xiaojukeji.know.streaming.km.common.enums.config.ConfigGroupEnum;
import com.xiaojukeji.know.streaming.km.common.enums.health.HealthCheckDimensionEnum; import com.xiaojukeji.know.streaming.km.common.enums.health.HealthCheckDimensionEnum;
import com.xiaojukeji.know.streaming.km.core.service.health.checkresult.HealthCheckResultService; import com.xiaojukeji.know.streaming.km.core.service.health.checkresult.HealthCheckResultService;
import com.xiaojukeji.know.streaming.km.core.service.health.score.HealthScoreService; import com.xiaojukeji.know.streaming.km.core.service.health.state.HealthStateService;
import io.swagger.annotations.Api; import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation; import io.swagger.annotations.ApiOperation;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@@ -28,7 +28,7 @@ import java.util.List;
@RequestMapping(ApiPrefix.API_V3_PREFIX) @RequestMapping(ApiPrefix.API_V3_PREFIX)
public class KafkaHealthController { public class KafkaHealthController {
@Autowired @Autowired
private HealthScoreService healthScoreService; private HealthStateService healthStateService;
@Autowired @Autowired
private HealthCheckResultService healthCheckResultService; private HealthCheckResultService healthCheckResultService;
@@ -40,12 +40,11 @@ public class KafkaHealthController {
@RequestParam(required = false) Integer dimensionCode) { @RequestParam(required = false) Integer dimensionCode) {
HealthCheckDimensionEnum dimensionEnum = HealthCheckDimensionEnum.getByCode(dimensionCode); HealthCheckDimensionEnum dimensionEnum = HealthCheckDimensionEnum.getByCode(dimensionCode);
if (!dimensionEnum.equals(HealthCheckDimensionEnum.UNKNOWN)) { if (!dimensionEnum.equals(HealthCheckDimensionEnum.UNKNOWN)) {
return Result.buildSuc(HealthScoreVOConverter.convert2HealthScoreResultDetailVOList(healthScoreService.getDimensionHealthScoreResult(clusterPhyId, dimensionEnum), false)); return Result.buildSuc(HealthScoreVOConverter.convert2HealthScoreResultDetailVOList(healthStateService.getDimensionHealthResult(clusterPhyId, dimensionEnum)));
} }
return Result.buildSuc(HealthScoreVOConverter.convert2HealthScoreResultDetailVOList( return Result.buildSuc(HealthScoreVOConverter.convert2HealthScoreResultDetailVOList(
healthScoreService.getClusterHealthScoreResult(clusterPhyId), healthStateService.getClusterHealthResult(clusterPhyId)
true
)); ));
} }
@@ -56,7 +55,7 @@ public class KafkaHealthController {
@PathVariable Integer dimensionCode, @PathVariable Integer dimensionCode,
@PathVariable String resName) { @PathVariable String resName) {
return Result.buildSuc(HealthScoreVOConverter.convert2HealthScoreBaseResultVOList( return Result.buildSuc(HealthScoreVOConverter.convert2HealthScoreBaseResultVOList(
healthScoreService.getResHealthScoreResult(clusterPhyId, dimensionCode, resName) healthStateService.getResHealthResult(clusterPhyId, dimensionCode, resName)
)); ));
} }