|
|
@@ -0,0 +1,114 @@
|
|
|
+package cn.seecoder.paas.service.scheduler;
|
|
|
+
|
|
|
+import cn.seecoder.paas.data.dao.ApplicationDAO;
|
|
|
+import cn.seecoder.paas.data.dao.EnvironmentDAO;
|
|
|
+import cn.seecoder.paas.data.dao.MetricsDAO;
|
|
|
+import cn.seecoder.paas.data.entity.Application;
|
|
|
+import cn.seecoder.paas.data.entity.Environment;
|
|
|
+import cn.seecoder.paas.service.MetricsService;
|
|
|
+import cn.seecoder.paas.service.facade.k8s.PodApi;
|
|
|
+import cn.seecoder.paas.service.facade.k8s.model.K8sObjectRequest;
|
|
|
+import cn.seecoder.paas.service.facade.k8s.model.Pod;
|
|
|
+import cn.seecoder.paas.service.model.converter.MetricsConverter;
|
|
|
+import cn.seecoder.paas.service.model.vo.MetricsVO;
|
|
|
+import cn.seecoder.paas.util.ApplicationProperties;
|
|
|
+import cn.seecoder.paas.util.enums.MetricsSource;
|
|
|
+import cn.seecoder.paas.util.enums.ResourceLabel;
|
|
|
+import io.fabric8.kubernetes.api.model.metrics.v1beta1.NodeMetricsList;
|
|
|
+import io.fabric8.kubernetes.api.model.metrics.v1beta1.PodMetrics;
|
|
|
+import io.fabric8.kubernetes.api.model.metrics.v1beta1.PodMetricsList;
|
|
|
+import io.fabric8.kubernetes.client.DefaultKubernetesClient;
|
|
|
+import org.springframework.beans.factory.annotation.Autowired;
|
|
|
+import org.springframework.scheduling.annotation.Scheduled;
|
|
|
+import org.springframework.stereotype.Component;
|
|
|
+
|
|
|
+import java.util.ArrayList;
|
|
|
+import java.util.Collections;
|
|
|
+import java.util.List;
|
|
|
+import java.util.stream.Collectors;
|
|
|
+
|
|
|
+@Component
|
|
|
+public class MetricsScheduler {
|
|
|
+
|
|
|
+ private final DefaultKubernetesClient fabricK8sClient;
|
|
|
+
|
|
|
+ private final ApplicationProperties applicationProperties;
|
|
|
+
|
|
|
+ private final EnvironmentDAO environmentDAO;
|
|
|
+
|
|
|
+ private final ApplicationDAO applicationDAO;
|
|
|
+
|
|
|
+ private final MetricsDAO metricsDAO;
|
|
|
+
|
|
|
+ private final MetricsService metricsService;
|
|
|
+
|
|
|
+ private final PodApi podApi;
|
|
|
+
|
|
|
+ @Autowired
|
|
|
+ public MetricsScheduler(DefaultKubernetesClient fabricK8sClient, ApplicationProperties applicationProperties, EnvironmentDAO environmentDAO, ApplicationDAO applicationDAO, MetricsDAO metricsDAO, MetricsService metricsService, PodApi podApi) {
|
|
|
+ this.fabricK8sClient = fabricK8sClient;
|
|
|
+ this.applicationProperties = applicationProperties;
|
|
|
+ this.environmentDAO = environmentDAO;
|
|
|
+ this.applicationDAO = applicationDAO;
|
|
|
+ this.metricsDAO = metricsDAO;
|
|
|
+ this.metricsService = metricsService;
|
|
|
+ this.podApi = podApi;
|
|
|
+ }
|
|
|
+
|
|
|
+ // 30秒采集一次
|
|
|
+ @Scheduled(cron = "30 * * * * ?")
|
|
|
+ public void statistic() {
|
|
|
+ String deploymentNamespace = applicationProperties.getDeploymentNamespace();
|
|
|
+ PodMetricsList podMetricsList = fabricK8sClient.top().pods().metrics(deploymentNamespace);
|
|
|
+ NodeMetricsList nodeMetricsList = fabricK8sClient.top().nodes().metrics();
|
|
|
+ // 集群
|
|
|
+ List<MetricsVO> metrics = metricsService.getMetricsSourceCurrentStatus(MetricsSource.CLUSTER, MetricsSource.CLUSTER.getCode());
|
|
|
+ metricsDAO.saveAll(metrics.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
|
|
|
+ // 节点
|
|
|
+ nodeMetricsList.getItems().forEach(nodeMetrics -> {
|
|
|
+ List<MetricsVO> values = metricsService.statNodeMetrics(Collections.singletonList(nodeMetrics), MetricsSource.NODE, nodeMetrics.getMetadata().getName());
|
|
|
+ if (!values.isEmpty())
|
|
|
+ metricsDAO.saveAll(values.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
|
|
|
+ });
|
|
|
+ // POD
|
|
|
+ podMetricsList.getItems().forEach(podMetrics -> {
|
|
|
+ List<MetricsVO> values = metricsService.statPodMetrics(Collections.singletonList(podMetrics), MetricsSource.POD, podMetrics.getMetadata().getName());
|
|
|
+ if (!values.isEmpty())
|
|
|
+ metricsDAO.saveAll(values.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
|
|
|
+ });
|
|
|
+ // ENVIRONMENT
|
|
|
+ List<Environment> environments = environmentDAO.findAll();
|
|
|
+ environments.forEach(environment -> {
|
|
|
+ String labelKey = ResourceLabel.RESOURCE.getCode();
|
|
|
+ String labelValue = ResourceLabel.RESOURCE.getGenerator().gen("environment", String.valueOf(environment.getAppId()), String.valueOf(environment.getId()));
|
|
|
+ K8sObjectRequest request = K8sObjectRequest.builder().namespace(applicationProperties.getDeploymentNamespace()).labels(Collections.singletonMap(labelKey, labelValue)).build();
|
|
|
+ List<Pod> pods = podApi.getByCondition(request);
|
|
|
+ List<String> podNames = pods.stream().map(item -> item.getName()).collect(Collectors.toList());
|
|
|
+ if (!pods.isEmpty()) {
|
|
|
+ List<PodMetrics> podMetricsValues = podMetricsList.getItems().stream().filter(item -> podNames.contains(item.getMetadata().getName())).collect(Collectors.toList());
|
|
|
+ List<MetricsVO> values = metricsService.statPodMetrics(podMetricsValues, MetricsSource.ENVIRONMENT, String.valueOf(environment.getId()));
|
|
|
+ if (!values.isEmpty())
|
|
|
+ metricsDAO.saveAll(values.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
|
|
|
+ }
|
|
|
+ });
|
|
|
+ // APPLICATION
|
|
|
+ List<Application> applications = applicationDAO.findAll();
|
|
|
+ applications.forEach(application -> {
|
|
|
+ List<Environment> appEnvs = environmentDAO.findAllByAppId(application.getId());
|
|
|
+ List<Pod> pods = new ArrayList<>();
|
|
|
+ appEnvs.forEach(environment -> {
|
|
|
+ String labelKey = ResourceLabel.RESOURCE.getCode();
|
|
|
+ String labelValue = ResourceLabel.RESOURCE.getGenerator().gen("environment", String.valueOf(environment.getAppId()), String.valueOf(environment.getId()));
|
|
|
+ K8sObjectRequest request = K8sObjectRequest.builder().namespace(applicationProperties.getDeploymentNamespace()).labels(Collections.singletonMap(labelKey, labelValue)).build();
|
|
|
+ pods.addAll(podApi.getByCondition(request));
|
|
|
+ });
|
|
|
+ List<String> podNames = pods.stream().map(item -> item.getName()).collect(Collectors.toList());
|
|
|
+ if (!pods.isEmpty()) {
|
|
|
+ List<PodMetrics> podMetricsValues = podMetricsList.getItems().stream().filter(item -> podNames.contains(item.getMetadata().getName())).collect(Collectors.toList());
|
|
|
+ List<MetricsVO> values = metricsService.statPodMetrics(podMetricsValues, MetricsSource.APPLICATION, String.valueOf(application.getId()));
|
|
|
+ if (!values.isEmpty())
|
|
|
+ metricsDAO.saveAll(values.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
|
|
|
+ }
|
|
|
+ });
|
|
|
+ }
|
|
|
+}
|