Przeglądaj źródła

refactor: 采集的定时上报| 添加定时功能

zhanghaijun 3 lat temu
rodzic
commit
0e08babc74

+ 2 - 0
src/main/java/com/cec/plugins/PluginsApplication.java

@@ -4,8 +4,10 @@ import org.springframework.boot.Banner;
 import org.springframework.boot.SpringApplication;
 import org.springframework.boot.SpringApplication;
 import org.springframework.boot.autoconfigure.SpringBootApplication;
 import org.springframework.boot.autoconfigure.SpringBootApplication;
 import org.springframework.scheduling.annotation.EnableAsync;
 import org.springframework.scheduling.annotation.EnableAsync;
+import org.springframework.scheduling.annotation.EnableScheduling;
 
 
 @EnableAsync
 @EnableAsync
+@EnableScheduling
 @SpringBootApplication
 @SpringBootApplication
 public class PluginsApplication {
 public class PluginsApplication {
 
 

+ 2 - 4
src/main/java/com/cec/plugins/config/AsyncConfig.java

@@ -4,7 +4,6 @@ import lombok.extern.slf4j.Slf4j;
 import org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler;
 import org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler;
 import org.springframework.context.annotation.Configuration;
 import org.springframework.context.annotation.Configuration;
 import org.springframework.scheduling.annotation.AsyncConfigurer;
 import org.springframework.scheduling.annotation.AsyncConfigurer;
-import org.springframework.scheduling.annotation.EnableAsync;
 import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
 import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
 
 
 import java.lang.reflect.Method;
 import java.lang.reflect.Method;
@@ -17,7 +16,6 @@ import java.util.concurrent.ThreadPoolExecutor;
  * @description 异步执行配置信息
  * @description 异步执行配置信息
  */
  */
 @Slf4j
 @Slf4j
-@EnableAsync
 @Configuration
 @Configuration
 public class AsyncConfig implements AsyncConfigurer {
 public class AsyncConfig implements AsyncConfigurer {
     @Override
     @Override
@@ -25,8 +23,8 @@ public class AsyncConfig implements AsyncConfigurer {
         ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
         ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
         executor.setCorePoolSize(8); // 核心线程数
         executor.setCorePoolSize(8); // 核心线程数
         executor.setMaxPoolSize(32);  // 最大线程数
         executor.setMaxPoolSize(32);  // 最大线程数
-        executor.setQueueCapacity(1000); // 队列大小
-        executor.setKeepAliveSeconds(600); // 线程最大空闲时间
+        executor.setQueueCapacity(30); // 队列大小
+        executor.setKeepAliveSeconds(60); // 线程最大空闲时间
         // 指定用于新创建的线程名称的前缀。
         // 指定用于新创建的线程名称的前缀。
         executor.setThreadNamePrefix("async-cec-");
         executor.setThreadNamePrefix("async-cec-");
         executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
         executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());

+ 1 - 7
src/main/java/com/cec/plugins/config/RunnerRegister.java

@@ -20,19 +20,13 @@ import org.springframework.stereotype.Component;
 public class RunnerRegister implements ApplicationRunner {
 public class RunnerRegister implements ApplicationRunner {
 
 
     private final TaskService taskService;
     private final TaskService taskService;
-    private final ApplicationContext ctx;
 
 
     @Override
     @Override
     public void run(ApplicationArguments args) {
     public void run(ApplicationArguments args) {
-        // 是否是开发环境
-        Boolean isDev = Convert.toBool(ctx.getEnvironment().getProperty("plugin.dev"), Boolean.FALSE);
         // 初始化
         // 初始化
         taskService.serverInit();
         taskService.serverInit();
         // 注册
         // 注册
         taskService.register();
         taskService.register();
-        // 心跳
-        if (!isDev) {
-            taskService.heartbeat();
-        }
+
     }
     }
 }
 }

+ 2 - 0
src/main/java/com/cec/plugins/config/properties/ProjectProperties.java

@@ -32,6 +32,8 @@ public class ProjectProperties {
         private int retries = 3;
         private int retries = 3;
         /** 间隔时间(毫秒) */
         /** 间隔时间(毫秒) */
         private long delayMillis = 1000L;
         private long delayMillis = 1000L;
+        /** 采集数据上报 */
+        private String collectCron = "0/1 * * * * ?";
 
 
         public String getHttpUtl(String url) {
         public String getHttpUtl(String url) {
             return this.getUrl() + url;
             return this.getUrl() + url;

+ 3 - 34
src/main/java/com/cec/plugins/controller/CollectController.java

@@ -9,6 +9,7 @@ import com.cec.plugins.enumerate.EstimateEnum;
 import com.cec.plugins.service.RedisClientEvaluate;
 import com.cec.plugins.service.RedisClientEvaluate;
 import com.cec.plugins.service.RedisClientSource;
 import com.cec.plugins.service.RedisClientSource;
 import com.cec.plugins.service.TaskService;
 import com.cec.plugins.service.TaskService;
+import com.cec.plugins.service.task.TaskCollectService;
 import lombok.RequiredArgsConstructor;
 import lombok.RequiredArgsConstructor;
 import org.springframework.beans.BeanUtils;
 import org.springframework.beans.BeanUtils;
 import org.springframework.scheduling.annotation.Async;
 import org.springframework.scheduling.annotation.Async;
@@ -37,9 +38,8 @@ import java.util.stream.Collectors;
 @RequestMapping("/config/api/v1")
 @RequestMapping("/config/api/v1")
 public class CollectController {
 public class CollectController {
 
 
-    private final RedisClientSource redisClientSource;
-    private final RedisClientEvaluate redisClientEvaluate;
     private final TaskService taskService;
     private final TaskService taskService;
+    private final TaskCollectService taskCollectService;
 
 
     /** 数据采集配置下发接口 */
     /** 数据采集配置下发接口 */
     @PostMapping("/sendPolicy")
     @PostMapping("/sendPolicy")
@@ -51,7 +51,7 @@ public class CollectController {
         // 上报结果
         // 上报结果
         reported(policy);
         reported(policy);
         // 开始上报采集数据
         // 开始上报采集数据
-        reportPolicy(policy);
+        taskCollectService.startReported(policy);
         return ResultData.success();
         return ResultData.success();
     }
     }
 
 
@@ -67,35 +67,4 @@ public class CollectController {
         data.setReason("执行成功");
         data.setReason("执行成功");
         taskService.reportPlatform("/config/api/v1/receivetaskstatus", data);
         taskService.reportPlatform("/config/api/v1/receivetaskstatus", data);
     }
     }
-
-    /** 上报采集数据 */
-    @Async
-    public void reportPolicy(DtoPolicy policyDto) {
-        DtoPolicy.ReportPolicy report = new DtoPolicy.ReportPolicy();
-        // 拷贝父类(BaseData)的属性
-        BeanUtils.copyProperties(policyDto, report);
-        List<DtoPolicy.ReportPolicyConfig> configList = new ArrayList<>();
-        for (DtoPolicy.PolicyConfig config : policyDto.getPolicy()) {
-            DtoPolicy.ReportPolicyConfig reportItem = new DtoPolicy.ReportPolicyConfig();
-            BeanUtils.copyProperties(config, reportItem);
-            for (DtoPolicy.PolicyItem policy : config.getCollect()) {
-                DtoPolicy.ReportPolicyItem item = new DtoPolicy.ReportPolicyItem();
-                BeanUtils.copyProperties(policy, item);
-                item.setValue(redisClientSource.getCacheStr(item.getKey()));
-                EstimateEnum estimateEnum = EstimateEnum.findEstimateEnum(item.getKey());
-                item.setGroup(Optional.ofNullable(estimateEnum).map(EstimateEnum::getGroup).orElse(""));
-                item.setScore(StrUtil.emptyToDefault(redisClientEvaluate.getCacheStr(item.getKey()), "0"));
-                reportItem.addCollect(item);
-            }
-            // 计算平均值
-            Double avg = reportItem.getCollect().stream()
-                .map(DtoPolicy.ReportPolicyItem::getScore)
-                .map(item -> Convert.toDouble(item, 0D))
-                .collect(Collectors.averagingDouble(Double::doubleValue));
-            reportItem.setScore(Convert.toStr(avg));
-            configList.add(reportItem);
-        }
-        report.setData(configList);
-        taskService.reportPlatform("/data/api/v1/collectData", report);
-    }
 }
 }

+ 6 - 0
src/main/java/com/cec/plugins/controller/TaskController.java

@@ -7,6 +7,7 @@ import com.cec.plugins.domain.ResultData;
 import com.cec.plugins.domain.reported.ReportedTypeFactory;
 import com.cec.plugins.domain.reported.ReportedTypeFactory;
 import com.cec.plugins.enumerate.TaskEnum;
 import com.cec.plugins.enumerate.TaskEnum;
 import com.cec.plugins.service.TaskService;
 import com.cec.plugins.service.TaskService;
+import com.cec.plugins.service.task.TaskCollectService;
 import lombok.RequiredArgsConstructor;
 import lombok.RequiredArgsConstructor;
 import org.springframework.validation.annotation.Validated;
 import org.springframework.validation.annotation.Validated;
 import org.springframework.web.bind.annotation.*;
 import org.springframework.web.bind.annotation.*;
@@ -27,6 +28,8 @@ public class TaskController {
 
 
     private final TaskService taskService;
     private final TaskService taskService;
 
 
+    private final TaskCollectService taskCollectService;
+
     /** 任务状态下发接口 */
     /** 任务状态下发接口 */
     @PostMapping("/taskStatus")
     @PostMapping("/taskStatus")
     public ResultData taskStatus(@RequestBody @Validated DtoTaskStatus dto) {
     public ResultData taskStatus(@RequestBody @Validated DtoTaskStatus dto) {
@@ -34,6 +37,9 @@ public class TaskController {
         if (!result) {
         if (!result) {
             return ResultData.fail("任务状态保存失败");
             return ResultData.fail("任务状态保存失败");
         }
         }
+        if (dto.getStatus() == 1) {
+            taskCollectService.stopReported();
+        }
         // 异步处理开始上报状态数据( 1、查询 2、上报 )
         // 异步处理开始上报状态数据( 1、查询 2、上报 )
         taskService.reported(TaskEnum.TaskInstanceId, ReportedTypeFactory.ReportedStatus.class, dto);
         taskService.reported(TaskEnum.TaskInstanceId, ReportedTypeFactory.ReportedStatus.class, dto);
         return ResultData.success();
         return ResultData.success();

+ 0 - 1
src/main/java/com/cec/plugins/service/TaskService.java

@@ -69,7 +69,6 @@ public class TaskService {
         }
         }
     }
     }
 
 
-    @Async
     public void heartbeat() {
     public void heartbeat() {
         ProjectProperties.PlatformModel platform = properties.getPlatform();
         ProjectProperties.PlatformModel platform = properties.getPlatform();
         String url = platform.getHttpUtl("/api/v1/serviceState");
         String url = platform.getHttpUtl("/api/v1/serviceState");

+ 89 - 0
src/main/java/com/cec/plugins/service/task/TaskCollectService.java

@@ -0,0 +1,89 @@
+package com.cec.plugins.service.task;
+
+import cn.hutool.core.convert.Convert;
+import cn.hutool.core.util.StrUtil;
+import com.cec.plugins.config.properties.ProjectProperties;
+import com.cec.plugins.domain.DtoPolicy;
+import com.cec.plugins.enumerate.EstimateEnum;
+import com.cec.plugins.service.RedisClientEvaluate;
+import com.cec.plugins.service.RedisClientSource;
+import com.cec.plugins.service.TaskService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.BeanUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.scheduling.annotation.Async;
+import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
+import org.springframework.scheduling.support.CronTrigger;
+import org.springframework.stereotype.Service;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.ScheduledFuture;
+import java.util.stream.Collectors;
+
+/**
+ * @author zhang
+ * @date 2023/8/31 11:27
+ * @description 采集数据定时上报
+ */
+@Slf4j
+@Service
+public class TaskCollectService {
+
+    @Autowired
+    private ProjectProperties properties;
+    @Autowired
+    private RedisClientSource redisClientSource;
+    @Autowired
+    private RedisClientEvaluate redisClientEvaluate;
+    @Autowired
+    private TaskService taskService;
+    @Autowired
+    private ThreadPoolTaskScheduler threadPoolTaskScheduler;
+    private ScheduledFuture<?> future;
+
+    public void startReported(DtoPolicy policyDto) {
+        stopReported();
+        String cron = properties.getPlatform().getCollectCron();
+        future = threadPoolTaskScheduler.schedule(() -> {
+            log.info("###### 执行上报采集数据 ######");
+            reportPolicy(policyDto);
+        }, triggerContext -> new CronTrigger(cron).nextExecutionTime(triggerContext));
+        log.info("###### 开始上报采集数据定时,cron: {} ######", cron);
+    }
+
+    public void stopReported() {
+        if (future != null) {
+            future.cancel(true);
+        }
+        log.info("###### 停止上报采集数据定时 ######");
+    }
+
+    @Async
+    public void reportPolicy(DtoPolicy policyDto) {
+        DtoPolicy.ReportPolicy report = new DtoPolicy.ReportPolicy();
+        // 拷贝父类(BaseData)的属性
+        BeanUtils.copyProperties(policyDto, report);
+        List<DtoPolicy.ReportPolicyConfig> configList = new ArrayList<>();
+        for (DtoPolicy.PolicyConfig config : policyDto.getPolicy()) {
+            DtoPolicy.ReportPolicyConfig reportItem = new DtoPolicy.ReportPolicyConfig();
+            BeanUtils.copyProperties(config, reportItem);
+            for (DtoPolicy.PolicyItem policy : config.getCollect()) {
+                DtoPolicy.ReportPolicyItem item = new DtoPolicy.ReportPolicyItem();
+                BeanUtils.copyProperties(policy, item);
+                item.setValue(redisClientSource.getCacheStr(item.getKey()));
+                EstimateEnum estimateEnum = EstimateEnum.findEstimateEnum(item.getKey());
+                item.setGroup(Optional.ofNullable(estimateEnum).map(EstimateEnum::getGroup).orElse(""));
+                item.setScore(StrUtil.emptyToDefault(redisClientEvaluate.getCacheStr(item.getKey()), "0"));
+                reportItem.addCollect(item);
+            }
+            // 计算平均值
+            Double avg = reportItem.getCollect().stream().map(DtoPolicy.ReportPolicyItem::getScore).map(item -> Convert.toDouble(item, 0D)).collect(Collectors.averagingDouble(Double::doubleValue));
+            reportItem.setScore(Convert.toStr(avg));
+            configList.add(reportItem);
+        }
+        report.setData(configList);
+        taskService.reportPlatform("/data/api/v1/collectData", report);
+    }
+}

+ 34 - 0
src/main/java/com/cec/plugins/service/task/TaskScheduleService.java

@@ -0,0 +1,34 @@
+package com.cec.plugins.service.task;
+
+import cn.hutool.core.convert.Convert;
+import com.cec.plugins.service.TaskService;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.context.ApplicationContext;
+import org.springframework.scheduling.annotation.Async;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+/**
+ * @author zhang
+ * @date 2023/8/31 11:17
+ * @description 多线程定时任务
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+public class TaskScheduleService {
+
+    private final TaskService taskService;
+    private final ApplicationContext ctx;
+
+    @Async
+    @Scheduled(cron = "*/5 * * * * ?")
+    public void heartbeat() {
+        // 是否是开发环境
+        Boolean isDev = Convert.toBool(ctx.getEnvironment().getProperty("plugin.dev"), Boolean.FALSE);
+        if (!isDev) {
+            taskService.heartbeat();
+        }
+    }
+}

+ 42 - 0
src/main/resources/application-prod.yml

@@ -0,0 +1,42 @@
+spring:
+  application:
+    name: plugins
+
+server:
+  port: 8080
+
+logging:
+  level:
+    org.springframework: warn
+    org.hibernate.validator: warn
+    org.apache: warn
+    com.cec.plugins.core.retry.RetryTemplate: warn
+    com.cec: info
+  config: classpath:logback.xml
+
+redisson:
+  source:
+    client-name: ${spring.application.name}
+    host: 172.16.20.1
+    port: 6379
+    password:
+    database: 0
+  evaluate:
+    client-name: ${spring.application.name}
+    host: 172.16.20.1
+    port: 3679
+    password:
+    database: 1
+
+project:
+  platform:
+    url: http://192.168.30.252:8080
+    retries: 3
+    delay-millis: 1000
+    collect-cron: 0/1 * * * * ?
+  plugin:
+    ip-address: 172.16.20.1
+    port: ${server.port}
+  simulation:
+    retries: 30
+    delay-millis: 500

+ 1 - 0
src/main/resources/application.yml

@@ -33,6 +33,7 @@ project:
     url: ${PLATFORM_URL:http://127.0.0.1:${server.port}}
     url: ${PLATFORM_URL:http://127.0.0.1:${server.port}}
     retries: ${PLATFORM_RETRIES:3}
     retries: ${PLATFORM_RETRIES:3}
     delay-millis: ${PLATFORM_DELAY-MILLIS:1000}
     delay-millis: ${PLATFORM_DELAY-MILLIS:1000}
+    collect-cron: 0/1 * * * * ?
   plugin:
   plugin:
     ip-address: ${PLUGIN_HOST:127.0.0.1}
     ip-address: ${PLUGIN_HOST:127.0.0.1}
     port: ${PLUGIN_PORT:${server.port}}
     port: ${PLUGIN_PORT:${server.port}}

+ 0 - 2
web/types/auto-imports.d.ts

@@ -305,7 +305,6 @@ import { UnwrapRef } from 'vue'
 declare module 'vue' {
 declare module 'vue' {
   interface ComponentCustomProperties {
   interface ComponentCustomProperties {
     readonly EffectScope: UnwrapRef<typeof import('vue')['EffectScope']>
     readonly EffectScope: UnwrapRef<typeof import('vue')['EffectScope']>
-    readonly ElMessage: UnwrapRef<typeof import('element-plus/es')['ElMessage']>
     readonly acceptHMRUpdate: UnwrapRef<typeof import('pinia')['acceptHMRUpdate']>
     readonly acceptHMRUpdate: UnwrapRef<typeof import('pinia')['acceptHMRUpdate']>
     readonly asyncComputed: UnwrapRef<typeof import('@vueuse/core')['asyncComputed']>
     readonly asyncComputed: UnwrapRef<typeof import('@vueuse/core')['asyncComputed']>
     readonly autoResetRef: UnwrapRef<typeof import('@vueuse/core')['autoResetRef']>
     readonly autoResetRef: UnwrapRef<typeof import('@vueuse/core')['autoResetRef']>
@@ -598,7 +597,6 @@ declare module 'vue' {
 declare module '@vue/runtime-core' {
 declare module '@vue/runtime-core' {
   interface ComponentCustomProperties {
   interface ComponentCustomProperties {
     readonly EffectScope: UnwrapRef<typeof import('vue')['EffectScope']>
     readonly EffectScope: UnwrapRef<typeof import('vue')['EffectScope']>
-    readonly ElMessage: UnwrapRef<typeof import('element-plus/es')['ElMessage']>
     readonly acceptHMRUpdate: UnwrapRef<typeof import('pinia')['acceptHMRUpdate']>
     readonly acceptHMRUpdate: UnwrapRef<typeof import('pinia')['acceptHMRUpdate']>
     readonly asyncComputed: UnwrapRef<typeof import('@vueuse/core')['asyncComputed']>
     readonly asyncComputed: UnwrapRef<typeof import('@vueuse/core')['asyncComputed']>
     readonly autoResetRef: UnwrapRef<typeof import('@vueuse/core')['autoResetRef']>
     readonly autoResetRef: UnwrapRef<typeof import('@vueuse/core')['autoResetRef']>