1
0

MetricsScheduler.java 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121
  1. package cn.seecoder.paas.service.scheduler;
  2. import cn.seecoder.paas.data.dao.ApplicationDAO;
  3. import cn.seecoder.paas.data.dao.EnvironmentDAO;
  4. import cn.seecoder.paas.data.dao.MetricsDAO;
  5. import cn.seecoder.paas.data.entity.Application;
  6. import cn.seecoder.paas.data.entity.Environment;
  7. import cn.seecoder.paas.service.MetricsService;
  8. import cn.seecoder.paas.service.facade.k8s.PodApi;
  9. import cn.seecoder.paas.service.facade.k8s.model.K8sObjectRequest;
  10. import cn.seecoder.paas.service.facade.k8s.model.Pod;
  11. import cn.seecoder.paas.service.model.converter.MetricsConverter;
  12. import cn.seecoder.paas.service.model.vo.MetricsVO;
  13. import cn.seecoder.paas.util.ApplicationProperties;
  14. import cn.seecoder.paas.util.enums.MetricsSource;
  15. import cn.seecoder.paas.util.enums.ResourceLabel;
  16. import io.fabric8.kubernetes.api.model.metrics.v1beta1.NodeMetricsList;
  17. import io.fabric8.kubernetes.api.model.metrics.v1beta1.PodMetrics;
  18. import io.fabric8.kubernetes.api.model.metrics.v1beta1.PodMetricsList;
  19. import io.fabric8.kubernetes.client.DefaultKubernetesClient;
  20. import org.springframework.beans.factory.annotation.Autowired;
  21. import org.springframework.scheduling.annotation.Scheduled;
  22. import org.springframework.stereotype.Component;
  23. import java.time.LocalDateTime;
  24. import java.util.ArrayList;
  25. import java.util.Collections;
  26. import java.util.List;
  27. import java.util.stream.Collectors;
  28. @Component
  29. public class MetricsScheduler {
  30. private final DefaultKubernetesClient fabricK8sClient;
  31. private final ApplicationProperties applicationProperties;
  32. private final EnvironmentDAO environmentDAO;
  33. private final ApplicationDAO applicationDAO;
  34. private final MetricsDAO metricsDAO;
  35. private final MetricsService metricsService;
  36. private final PodApi podApi;
  37. @Autowired
  38. public MetricsScheduler(DefaultKubernetesClient fabricK8sClient, ApplicationProperties applicationProperties, EnvironmentDAO environmentDAO, ApplicationDAO applicationDAO, MetricsDAO metricsDAO, MetricsService metricsService, PodApi podApi) {
  39. this.fabricK8sClient = fabricK8sClient;
  40. this.applicationProperties = applicationProperties;
  41. this.environmentDAO = environmentDAO;
  42. this.applicationDAO = applicationDAO;
  43. this.metricsDAO = metricsDAO;
  44. this.metricsService = metricsService;
  45. this.podApi = podApi;
  46. }
  47. // 30秒采集一次
  48. @Scheduled(cron = "30 * * * * ?")
  49. public void statistic() {
  50. String deploymentNamespace = applicationProperties.getDeploymentNamespace();
  51. PodMetricsList podMetricsList = fabricK8sClient.top().pods().metrics(deploymentNamespace);
  52. NodeMetricsList nodeMetricsList = fabricK8sClient.top().nodes().metrics();
  53. // 集群
  54. List<MetricsVO> metrics = metricsService.getMetricsSourceCurrentStatus(MetricsSource.CLUSTER, MetricsSource.CLUSTER.getCode());
  55. metricsDAO.saveAll(metrics.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
  56. // 节点
  57. nodeMetricsList.getItems().forEach(nodeMetrics -> {
  58. List<MetricsVO> values = metricsService.statNodeMetrics(Collections.singletonList(nodeMetrics), MetricsSource.NODE, nodeMetrics.getMetadata().getName());
  59. if (!values.isEmpty())
  60. metricsDAO.saveAll(values.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
  61. });
  62. // POD
  63. podMetricsList.getItems().forEach(podMetrics -> {
  64. List<MetricsVO> values = metricsService.statPodMetrics(Collections.singletonList(podMetrics), MetricsSource.POD, podMetrics.getMetadata().getName());
  65. if (!values.isEmpty())
  66. metricsDAO.saveAll(values.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
  67. });
  68. // ENVIRONMENT
  69. List<Environment> environments = environmentDAO.findAll();
  70. environments.forEach(environment -> {
  71. String labelKey = ResourceLabel.RESOURCE.getCode();
  72. String labelValue = ResourceLabel.RESOURCE.getGenerator().gen("environment", String.valueOf(environment.getAppId()), String.valueOf(environment.getId()));
  73. K8sObjectRequest request = K8sObjectRequest.builder().namespace(applicationProperties.getDeploymentNamespace()).labels(Collections.singletonMap(labelKey, labelValue)).build();
  74. List<Pod> pods = podApi.getByCondition(request);
  75. List<String> podNames = pods.stream().map(item -> item.getName()).collect(Collectors.toList());
  76. if (!pods.isEmpty()) {
  77. List<PodMetrics> podMetricsValues = podMetricsList.getItems().stream().filter(item -> podNames.contains(item.getMetadata().getName())).collect(Collectors.toList());
  78. List<MetricsVO> values = metricsService.statPodMetrics(podMetricsValues, MetricsSource.ENVIRONMENT, String.valueOf(environment.getId()));
  79. if (!values.isEmpty())
  80. metricsDAO.saveAll(values.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
  81. }
  82. });
  83. // APPLICATION
  84. List<Application> applications = applicationDAO.findAll();
  85. applications.forEach(application -> {
  86. List<Environment> appEnvs = environmentDAO.findAllByAppId(application.getId());
  87. List<Pod> pods = new ArrayList<>();
  88. appEnvs.forEach(environment -> {
  89. String labelKey = ResourceLabel.RESOURCE.getCode();
  90. String labelValue = ResourceLabel.RESOURCE.getGenerator().gen("environment", String.valueOf(environment.getAppId()), String.valueOf(environment.getId()));
  91. K8sObjectRequest request = K8sObjectRequest.builder().namespace(applicationProperties.getDeploymentNamespace()).labels(Collections.singletonMap(labelKey, labelValue)).build();
  92. pods.addAll(podApi.getByCondition(request));
  93. });
  94. List<String> podNames = pods.stream().map(item -> item.getName()).collect(Collectors.toList());
  95. if (!pods.isEmpty()) {
  96. List<PodMetrics> podMetricsValues = podMetricsList.getItems().stream().filter(item -> podNames.contains(item.getMetadata().getName())).collect(Collectors.toList());
  97. List<MetricsVO> values = metricsService.statPodMetrics(podMetricsValues, MetricsSource.APPLICATION, String.valueOf(application.getId()));
  98. if (!values.isEmpty())
  99. metricsDAO.saveAll(values.stream().map(MetricsConverter::convertToEntity).collect(Collectors.toList()));
  100. }
  101. });
  102. }
  103. @Scheduled(cron = "0 0 4 * * ?") // 每天凌晨4点执行
  104. public void cleanupOldMetrics() {
  105. LocalDateTime twoDaysAgo = LocalDateTime.now().minusDays(2);
  106. metricsDAO.deleteByMetricsTimeBefore(twoDaysAgo);
  107. }
  108. }