| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121 |
- 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.time.LocalDateTime;
- 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()));
- }
- });
- }
- @Scheduled(cron = "0 0 4 * * ?") // 每天凌晨4点执行
- public void cleanupOldMetrics() {
- LocalDateTime twoDaysAgo = LocalDateTime.now().minusDays(2);
- metricsDAO.deleteByMetricsTimeBefore(twoDaysAgo);
- }
- }
|