diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/TransportTaskTopEnum.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/TaskTopEnum.java similarity index 80% rename from jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/TransportTaskTopEnum.java rename to jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/TaskTopEnum.java index 42f6b3d..0e20507 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/TransportTaskTopEnum.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/TaskTopEnum.java @@ -3,7 +3,7 @@ package org.jeecg.common.constant.enums; /** * 输运模拟任务置顶说明枚举 */ -public enum TransportTaskTopEnum { +public enum TaskTopEnum { /** * 未置顶 @@ -17,7 +17,7 @@ public enum TransportTaskTopEnum { private Integer value; - TransportTaskTopEnum(Integer value) { + TaskTopEnum(Integer value) { this.value = value; } diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/SourceRebuildTask.java b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/SourceRebuildTask.java index 8907430..3ace1ff 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/SourceRebuildTask.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/SourceRebuildTask.java @@ -86,6 +86,7 @@ public class SourceRebuildTask implements Serializable { */ @NotNull(message = "srs时间周期-开始日期不能为空",groups = {InsertGroup.class, UpdateGroup.class}) @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd") + @DateTimeFormat(pattern = "yyyy-MM-dd") @TableField(value = "srs_start_time") private LocalDate srsStartTime; @@ -94,6 +95,7 @@ public class SourceRebuildTask implements Serializable { */ @NotNull(message = "srs时间周期-结束日期不能为空",groups = {InsertGroup.class, UpdateGroup.class}) @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd") + @DateTimeFormat(pattern = "yyyy-MM-dd") @TableField(value = "srs_end_time") private LocalDate srsEndTime; @@ -127,14 +129,16 @@ public class SourceRebuildTask implements Serializable { * 释放源开始释放时间 */ @TableField(value = "release_start_time") - @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd") + @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss") + @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss") private LocalDateTime releaseStartTime; /** * 释放源结束释放时间 */ @TableField(value = "release_end_time") - @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd") + @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss") + @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss") private LocalDateTime releaseEndTime; /** diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/consumer/RebuildTaskConsumerHandler.java b/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/consumer/RebuildTaskConsumerHandler.java index fcb246b..b42bf69 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/consumer/RebuildTaskConsumerHandler.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/consumer/RebuildTaskConsumerHandler.java @@ -39,7 +39,7 @@ public class RebuildTaskConsumerHandler { MessageConsumerThread messageConsumerThread = new MessageConsumerThread(); messageConsumerThread.setName("rebuild-task-thread"); messageConsumerThread.start(); - log.info("启动源项重建任务消费线程----------------------"); + log.info("启动源项重建任务消费线程"); } /** diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/service/SourceRebuildTaskService.java b/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/service/SourceRebuildTaskService.java index 30e79d9..507dc2a 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/service/SourceRebuildTaskService.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/service/SourceRebuildTaskService.java @@ -28,4 +28,10 @@ public interface SourceRebuildTaskService extends IService { * @param task */ void setInpectionFailed(SourceRebuildTask task); + + /** + * 取消置顶 + * @param taskId + */ + void setCancelTaskTop(Integer taskId); } diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/service/impl/SourceRebuildTaskServiceImpl.java b/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/service/impl/SourceRebuildTaskServiceImpl.java index e06f850..9b7d207 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/service/impl/SourceRebuildTaskServiceImpl.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/service/impl/SourceRebuildTaskServiceImpl.java @@ -4,6 +4,7 @@ import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import lombok.RequiredArgsConstructor; import org.jeecg.common.constant.CommonConstant; import org.jeecg.common.constant.enums.SourceRebuildTaskStatusEnum; +import org.jeecg.common.constant.enums.TaskTopEnum; import org.jeecg.common.properties.ServerProperties; import org.jeecg.common.util.RedisUtil; import org.jeecg.modules.base.entity.SourceRebuildTask; @@ -11,6 +12,7 @@ import org.jeecg.modules.base.mapper.SourceRebuildTaskMapper; import org.jeecg.rebuild.service.SourceRebuildTaskService; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; +import java.util.Objects; /** * 源项重建任务 @@ -62,4 +64,21 @@ public class SourceRebuildTaskServiceImpl extends ServiceImpl lines = new ArrayList<>(); lines.add(title); @@ -244,8 +233,7 @@ public class SourceRebuildTaskExec extends Thread{ String srsDataLog = "----------------------------------------生成SRS数据输入文件----------------------------------------"; this.generateLog(srsDataLog); this.srsFilesPath.forEach(this::generateLog); -// String inputPath = sourceRebuildProperties.getRInput()+File.separator+ "srsfilelist_subexp1.dat"; - String inputPath = sourceRebuildProperties.getRInput()+"/"+ "srsfilelist_subexp1.dat"; + String inputPath = sourceRebuildProperties.getRInput()+File.separator+ "srsfilelist_subexp1.dat"; jschRemoteRunner.writeFile(inputPath,String.join("",this.srsFilesPath)); } @@ -330,7 +318,6 @@ public class SourceRebuildTaskExec extends Thread{ " )" + " })" + "})"; -// REXP outputREXP = conn.eval("capture.output({source(\"R_scripts/main.R\")})"); REXP outputREXP = conn.eval(rCmd); // 将输出作为字符串数组获取 String[] outputLines = outputREXP.asStrings(); @@ -370,8 +357,7 @@ public class SourceRebuildTaskExec extends Thread{ private String getOutputPath(){ StringBuilder outputPath = new StringBuilder(); outputPath.append(sourceRebuildProperties.getROutput()); - outputPath.append("/"); -// outputPath.append(File.separator); + outputPath.append(File.separator); outputPath.append(sourceRebuildTask.getTaskName()); //创建任务输出目录 this.jschRemoteRunner.mkdir(outputPath.toString()); diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/transport/consumer/TranTaskConsumerHandler.java b/jeecg-model-consumer/src/main/java/org/jeecg/transport/consumer/TranTaskConsumerHandler.java index 70b0de8..b02b6ef 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/transport/consumer/TranTaskConsumerHandler.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/transport/consumer/TranTaskConsumerHandler.java @@ -1,23 +1,15 @@ package org.jeecg.transport.consumer; -import cn.hutool.core.collection.CollUtil; -import com.alibaba.fastjson2.JSON; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.apache.rocketmq.client.consumer.DefaultLitePullConsumer; -import org.apache.rocketmq.client.exception.MQClientException; -import org.apache.rocketmq.common.message.MessageExt; import org.jeecg.common.constant.CommonConstant; -import org.jeecg.common.constant.RocketMQTopConstant; import org.jeecg.common.constant.enums.TransportTaskStatusEnum; -import org.jeecg.common.constant.enums.TransportTaskTopEnum; import org.jeecg.common.properties.DataFusionProperties; import org.jeecg.common.properties.ServerProperties; import org.jeecg.common.properties.SystemStorageProperties; import org.jeecg.common.properties.TransportSimulationProperties; import org.jeecg.common.util.RedisUtil; -import org.jeecg.modules.base.dto.TransportTaskDTO; import org.jeecg.modules.base.entity.TransportTask; import org.jeecg.modules.base.mapper.*; import org.jeecg.transport.consumer.china.AbstractTaskMsgHandler; @@ -25,10 +17,8 @@ import org.jeecg.transport.consumer.china.Server11TaskHandler; import org.jeecg.transport.service.StationDataService; import org.jeecg.transport.service.StationsModValService; import org.jeecg.transport.service.TransportTaskService; -import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; -import java.util.List; import java.util.Objects; import java.util.concurrent.TimeUnit; @@ -59,7 +49,7 @@ public class TranTaskConsumerHandler{ MessageConsumerThread messageConsumerThread = new MessageConsumerThread(); messageConsumerThread.setName("transport-task-thread"); messageConsumerThread.start(); - log.info("启动输运模拟任务消费线程----------------------"); + log.info("启动输运模拟任务消费线程"); } /** diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/BackwardTaskExec.java b/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/BackwardTaskExec.java index 84f82db..910ceed 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/BackwardTaskExec.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/BackwardTaskExec.java @@ -89,6 +89,8 @@ public class BackwardTaskExec extends AbstractTaskExec { super.setTaskRunFlag(); //修改任务状态为执行中 super.transportTaskService.updateTaskStatus(super.transportTask.getId(), TransportTaskStatusEnum.IN_OPERATION.getValue()); + //任务开始运行后,如果之前是置顶状态,则设置任务取消置顶 + this.transportTaskService.setCancelTaskTop(this.transportTask.getId()); //如果此任务已存在历史日志,先清除 super.transportTaskService.deleteTaskLog(super.transportTask.getId()); //执行模拟 diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/ForwardTaskExec.java b/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/ForwardTaskExec.java index 1134de6..beff7c0 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/ForwardTaskExec.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/ForwardTaskExec.java @@ -105,6 +105,8 @@ public class ForwardTaskExec extends AbstractTaskExec { super.setTaskRunFlag(); //修改任务状态为执行中 super.transportTaskService.updateTaskStatus(super.transportTask.getId(), TransportTaskStatusEnum.IN_OPERATION.getValue()); + //任务开始运行后,如果之前是置顶状态,则设置任务取消置顶 + this.transportTaskService.setCancelTaskTop(this.transportTask.getId()); //如果此任务已存在历史日志,先清除 super.transportTaskService.deleteTaskLog(super.transportTask.getId()); //执行模拟 diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/transport/service/TransportTaskService.java b/jeecg-model-consumer/src/main/java/org/jeecg/transport/service/TransportTaskService.java index 58668bc..7119b46 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/transport/service/TransportTaskService.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/transport/service/TransportTaskService.java @@ -49,4 +49,10 @@ public interface TransportTaskService extends IService { * @param transportTask */ void setInpectionFailed(TransportTask transportTask); + + /** + * 取消置顶 + * @param taskId + */ + void setCancelTaskTop(Integer taskId); } diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/transport/service/impl/TransportTaskServiceImpl.java b/jeecg-model-consumer/src/main/java/org/jeecg/transport/service/impl/TransportTaskServiceImpl.java index df637cc..d62626a 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/transport/service/impl/TransportTaskServiceImpl.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/transport/service/impl/TransportTaskServiceImpl.java @@ -13,8 +13,6 @@ import org.jeecg.modules.base.mapper.*; import org.jeecg.transport.service.TransportTaskService; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; - -import java.time.LocalDateTime; import java.util.*; /** @@ -53,7 +51,7 @@ public class TransportTaskServiceImpl extends ServiceImpl { */ void setInpectionFailed(WeatherTask weatherTask); + /** + * 取消置顶 + * @param taskId + */ + void setCancelTaskTop(Integer taskId); + } diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/weather/service/impl/WeatherTaskServiceImpl.java b/jeecg-model-consumer/src/main/java/org/jeecg/weather/service/impl/WeatherTaskServiceImpl.java index fbfae12..e54b934 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/weather/service/impl/WeatherTaskServiceImpl.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/weather/service/impl/WeatherTaskServiceImpl.java @@ -4,6 +4,7 @@ import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import lombok.RequiredArgsConstructor; import org.jeecg.common.constant.CommonConstant; +import org.jeecg.common.constant.enums.TaskTopEnum; import org.jeecg.common.constant.enums.WeatherTaskStatusEnum; import org.jeecg.common.properties.ServerProperties; import org.jeecg.common.util.RedisUtil; @@ -14,6 +15,7 @@ import org.jeecg.modules.base.mapper.WeatherTaskMapper; import org.jeecg.weather.service.WeatherTaskService; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; +import java.util.Objects; /** * 天气预报预测任务管理 @@ -90,4 +92,22 @@ public class WeatherTaskServiceImpl extends ServiceImpl{ ProgressQueue.getInstance().offer(new ProgressEvent(this.weatherTask.getId(),log)); + if((log.startsWith(EXCEPTION_LOG1) || log.startsWith(EXCEPTION_LOG2)) && !this.failureFlag){ + this.failureFlag = true; + this.weatherTaskService.updateTaskStatus(this.weatherTask.getId(),WeatherTaskStatusEnum.FAILURE.getValue()); + } }) .doOnError(e->{ throw new RuntimeException(e); @@ -236,11 +266,16 @@ public class WeatherForecastTaskExec extends AbstractWeatherTask { ProgressQueue.getInstance().offer(new ProgressEvent(this.weatherTask.getId(),log)); return; } - //把预测好的及格式化后的气象文件移动到最终目录 + String formatFilesFinalPath = this.getFormatFilesFinalStoragePath(); String sourceFilesFinalPath = this.getSourceFilesFinalStoragePath(); - FileUtil.move(new File(gribCopyPath),new File(sourceFilesFinalPath),true); - FileUtil.move(new File(flexpartFormatPath),new File(formatFilesFinalPath),true); + try{ + //把预测好的及格式化后的气象文件移动到最终目录 + this.moveAll(gribCopyPath,sourceFilesFinalPath); + this.moveAll(flexpartFormatPath,formatFilesFinalPath); + }catch (Exception e){ + throw new RuntimeException("文件移动出现错误",e); + } //处理文件入库 List dataList = new ArrayList<>(); for(File sourceFile : sourceFiles){ @@ -252,8 +287,9 @@ public class WeatherForecastTaskExec extends AbstractWeatherTask { weatherData.setDataStartTime(Grib2TimeReader.readValidTime(sourceFilesFinalPath+File.separator+sourceFile.getName())); weatherData.setDataSource(weatherTask.getPredictionModel()); weatherData.setFilePath(sourceFilesFinalPath+File.separator+sourceFile.getName()); - weatherData.setFormatFilePath(formatFilesFinalPath+File.separator+sourceFile.getName()); + weatherData.setFormatFilePath(formatFilesFinalPath+File.separator+sourceFile.getName().substring(0,sourceFile.getName().lastIndexOf(".")+1)+ WeatherFileSuffixEnum.GRIB2.getValue()); weatherData.setTaskId(this.weatherTask.getId()); + weatherData.setMd5Value(this.getGribFileMD5(sourceFilesFinalPath+File.separator+sourceFile.getName())); dataList.add(weatherData); }catch (Exception e){ String logContent = "读取"+sourceFilesFinalPath+File.separator+sourceFile.getName()+"文件时间参数出现错误"; @@ -266,6 +302,36 @@ public class WeatherForecastTaskExec extends AbstractWeatherTask { } } + /** + * 移动所有文件到指定目录 + * @param srcDirPath + * @param destDirPath + */ + private void moveAll(String srcDirPath,String destDirPath) throws IOException { + File srcDir = new File(srcDirPath); + File destDir = new File(destDirPath); + if(!destDir.exists()){ + FileUtils.forceMkdir(destDir); + } + File[] files = srcDir.listFiles(); + if(ArrayUtils.isEmpty(files)){ + return; + } + for(File file : files){ + FileUtils.moveToDirectory(file,destDir,false); + } + } + + /** + * 获取GRIB文件的MD5唯一值 + */ + private String getGribFileMD5(String filePath) throws IOException { + try (FileInputStream fis = new FileInputStream(filePath)) { + // 底层自动采用流式读取,内存占用极低 + return DigestUtils.md5Hex(fis); + } + } + /** * 获取盘古模型请求命令 * @return @@ -307,7 +373,7 @@ public class WeatherForecastTaskExec extends AbstractWeatherTask { map.put("lead_time",this.weatherTask.getLeadTime()); map.put("class","od"); map.put("assets","assets-graphcast"); - map.put("workdir",systemStorageProperties.getPanguModelExecPath()); + map.put("workdir",systemStorageProperties.getGraphcastModelExecPath()); map.put("split_dir",getGribCopyPath()); map.put("path",buildOutputFilePath()); @@ -347,7 +413,7 @@ public class WeatherForecastTaskExec extends AbstractWeatherTask { if(WeatherDataSourceEnum.PANGU.getKey().equals(weatherTask.getPredictionModel())){ path.append(systemStorageProperties.getPanguModelExecPath()); }else if(WeatherDataSourceEnum.GRAPHCAST.getKey().equals(weatherTask.getPredictionModel())){ - path.append(systemStorageProperties.getPanguModelExecPath()); + path.append(systemStorageProperties.getGraphcastModelExecPath()); } path.append(File.separator); path.append(this.weatherTask.getId()); @@ -380,9 +446,9 @@ public class WeatherForecastTaskExec extends AbstractWeatherTask { */ private String getSourceFilesFinalStoragePath(){ if(WeatherDataSourceEnum.PANGU.getKey().equals(weatherTask.getPredictionModel())){ - return systemStorageProperties.getPanguModelExecPath()+"/source"; + return systemStorageProperties.getPanguDataPath()+"/source"; }else if(WeatherDataSourceEnum.GRAPHCAST.getKey().equals(weatherTask.getPredictionModel())){ - return systemStorageProperties.getGraphcastModelExecPath()+"/source"; + return systemStorageProperties.getGraphcastDataPath()+"/source"; } return ""; } @@ -393,9 +459,9 @@ public class WeatherForecastTaskExec extends AbstractWeatherTask { */ private String getFormatFilesFinalStoragePath(){ if(WeatherDataSourceEnum.PANGU.getKey().equals(weatherTask.getPredictionModel())){ - return systemStorageProperties.getPanguModelExecPath()+"/format"; + return systemStorageProperties.getPanguDataPath()+"/format"; }else if(WeatherDataSourceEnum.GRAPHCAST.getKey().equals(weatherTask.getPredictionModel())){ - return systemStorageProperties.getGraphcastModelExecPath()+"/format"; + return systemStorageProperties.getGraphcastDataPath()+"/format"; } return ""; } diff --git a/jeecg-module-data-analyze/src/main/java/org/jeecg/controller/DataAnalysisController.java b/jeecg-module-data-analyze/src/main/java/org/jeecg/controller/DataAnalysisController.java index 0a1bed4..2b4ae68 100644 --- a/jeecg-module-data-analyze/src/main/java/org/jeecg/controller/DataAnalysisController.java +++ b/jeecg-module-data-analyze/src/main/java/org/jeecg/controller/DataAnalysisController.java @@ -14,7 +14,6 @@ import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; - import java.util.Date; import java.util.List; diff --git a/jeecg-module-source-rebuild/src/main/java/org/jeecg/controller/SourceRebuildTaskController.java b/jeecg-module-source-rebuild/src/main/java/org/jeecg/controller/SourceRebuildTaskController.java index aac99e9..1402f96 100644 --- a/jeecg-module-source-rebuild/src/main/java/org/jeecg/controller/SourceRebuildTaskController.java +++ b/jeecg-module-source-rebuild/src/main/java/org/jeecg/controller/SourceRebuildTaskController.java @@ -87,4 +87,20 @@ public class SourceRebuildTaskController{ sourceRebuildTaskService.runTask(taskId); return Result.OK(); } + + @AutoLog(value = "设置任务置顶") + @Operation(summary = "设置任务置顶") + @PutMapping("setTaskTop") + public Result setTaskTop(@NotNull(message = "任务ID不能为空") Integer taskId){ + sourceRebuildTaskService.setTaskTop(taskId); + return Result.OK(); + } + + @AutoLog(value = "设置取消任务置顶") + @Operation(summary = "设置取消任务置顶") + @PutMapping("setCancelTaskTop") + public Result setCancelTaskTop(@NotNull(message = "任务ID不能为空") Integer taskId){ + sourceRebuildTaskService.setCancelTaskTop(taskId); + return Result.OK(); + } } diff --git a/jeecg-module-source-rebuild/src/main/java/org/jeecg/service/SourceRebuildTaskService.java b/jeecg-module-source-rebuild/src/main/java/org/jeecg/service/SourceRebuildTaskService.java index 3a76cbc..fc1e931 100644 --- a/jeecg-module-source-rebuild/src/main/java/org/jeecg/service/SourceRebuildTaskService.java +++ b/jeecg-module-source-rebuild/src/main/java/org/jeecg/service/SourceRebuildTaskService.java @@ -62,4 +62,16 @@ public interface SourceRebuildTaskService extends IService { * @return */ List getTaskLog(Integer taskId); + + /** + * 设置任务置顶 + * @param taskId + */ + void setTaskTop(Integer taskId); + + /** + * 取消置顶 + * @param taskId + */ + void setCancelTaskTop(Integer taskId); } diff --git a/jeecg-module-source-rebuild/src/main/java/org/jeecg/service/impl/SourceRebuildTaskServiceImpl.java b/jeecg-module-source-rebuild/src/main/java/org/jeecg/service/impl/SourceRebuildTaskServiceImpl.java index 775e305..d4c72c6 100644 --- a/jeecg-module-source-rebuild/src/main/java/org/jeecg/service/impl/SourceRebuildTaskServiceImpl.java +++ b/jeecg-module-source-rebuild/src/main/java/org/jeecg/service/impl/SourceRebuildTaskServiceImpl.java @@ -8,6 +8,7 @@ import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import lombok.RequiredArgsConstructor; import org.apache.commons.lang3.StringUtils; import org.jeecg.common.constant.enums.SourceRebuildTaskStatusEnum; +import org.jeecg.common.constant.enums.TaskTopEnum; import org.jeecg.common.properties.SourceRebuildProperties; import org.jeecg.common.system.query.PageRequest; import org.jeecg.modules.base.entity.SourceRebuildMonitoringData; @@ -77,8 +78,6 @@ public class SourceRebuildTaskServiceImpl extends ServiceImpl setCancelTaskTop(@NotNull(message = "任务ID不能为空") Integer taskId){ transportTaskService.setCancelTaskTop(taskId); diff --git a/jeecg-module-transport/src/main/java/org/jeecg/service/impl/TransportTaskServiceImpl.java b/jeecg-module-transport/src/main/java/org/jeecg/service/impl/TransportTaskServiceImpl.java index f69905b..672043f 100644 --- a/jeecg-module-transport/src/main/java/org/jeecg/service/impl/TransportTaskServiceImpl.java +++ b/jeecg-module-transport/src/main/java/org/jeecg/service/impl/TransportTaskServiceImpl.java @@ -107,7 +107,7 @@ public class TransportTaskServiceImpl extends ServiceImpl handleStaticDataToDB(String path,Integer dataSource){ - weatherDataService.handleStaticDataToDB(path,dataSource); - return Result.OK(); - } - @AutoLog(value = "关联气象数据") @Operation(summary = "关联气象数据") @PutMapping("linkedData") diff --git a/jeecg-module-weather/src/main/java/org/jeecg/controller/WeatherTaskController.java b/jeecg-module-weather/src/main/java/org/jeecg/controller/WeatherTaskController.java index 1b15270..2c70e7d 100644 --- a/jeecg-module-weather/src/main/java/org/jeecg/controller/WeatherTaskController.java +++ b/jeecg-module-weather/src/main/java/org/jeecg/controller/WeatherTaskController.java @@ -2,7 +2,6 @@ package org.jeecg.controller; import com.baomidou.mybatisplus.core.metadata.IPage; import io.swagger.v3.oas.annotations.Operation; -import jakarta.validation.constraints.NotBlank; import jakarta.validation.constraints.NotNull; import lombok.RequiredArgsConstructor; import org.jeecg.common.api.vo.Result; @@ -86,4 +85,21 @@ public class WeatherTaskController { public Result getTaskLog(@NotNull(message = "预测任务ID不能为空") Integer taskId){ return Result.OK(weatherTaskService.getTaskLog(taskId)); } + + @AutoLog(value = "设置任务置顶") + @Operation(summary = "设置任务置顶") + @PutMapping("setTaskTop") + public Result setTaskTop(@NotNull(message = "任务ID不能为空") Integer taskId){ + weatherTaskService.setTaskTop(taskId); + return Result.OK(); + } + + @AutoLog(value = "设置取消任务置顶") + @Operation(summary = "设置取消任务置顶") + @PutMapping("setCancelTaskTop") + public Result setCancelTaskTop(@NotNull(message = "任务ID不能为空") Integer taskId){ + weatherTaskService.setCancelTaskTop(taskId); + return Result.OK(); + } + } diff --git a/jeecg-module-weather/src/main/java/org/jeecg/job/DownloadT1hJob.java b/jeecg-module-weather/src/main/java/org/jeecg/job/DownloadT1hJob.java index abdd0f3..c923a1e 100644 --- a/jeecg-module-weather/src/main/java/org/jeecg/job/DownloadT1hJob.java +++ b/jeecg-module-weather/src/main/java/org/jeecg/job/DownloadT1hJob.java @@ -195,24 +195,42 @@ public class DownloadT1hJob { .filter(Files::isRegularFile) .filter(path -> path.toString().toLowerCase().endsWith(WeatherFileSuffixEnum.GRIB2.getValue())) .map(Path::toFile) - .map(file -> extractFileInfo(file, baseTime)) + .map(file -> { + try { + return extractFileInfo(file, baseTime); + } catch (IOException e) { + throw new RuntimeException(e); + } + }) .collect(Collectors.toList()); } catch (IOException e) { throw new RuntimeException("读取文件夹失败: " + folderPath, e); } } - private WeatherData extractFileInfo(File file, String baseTime) { + private WeatherData extractFileInfo(File file, String baseTime) throws IOException { + String gribFileMD5 = this.getGribFileMD5(file.getAbsolutePath()); WeatherData data = new WeatherData(); data.setFileName(file.getName()); data.setFileExt(getFileExtension(file.getName())); data.setFilePath(file.getAbsolutePath()); data.setDataSource(WeatherDataSourceEnum.T1H.getKey()); data.setDataStartTime(parseStartTimeFromFileName(file.getName())); + data.setMd5Value(gribFileMD5); data.setTimeBatch(baseTime); return data; } + /** + * 获取GRIB文件的MD5唯一值 + */ + private String getGribFileMD5(String filePath) throws IOException { + try (FileInputStream fis = new FileInputStream(filePath)) { + // 底层自动采用流式读取,内存占用极低 + return DigestUtils.md5Hex(fis); + } + } + private LocalDateTime parseStartTimeFromFileName(String fileName) { // 从文件名解析时间 示例:"T1H_20251029_00.nc" 这样的格式 try { diff --git a/jeecg-module-weather/src/main/java/org/jeecg/service/WeatherDataService.java b/jeecg-module-weather/src/main/java/org/jeecg/service/WeatherDataService.java index e1d41ce..8a3a2aa 100644 --- a/jeecg-module-weather/src/main/java/org/jeecg/service/WeatherDataService.java +++ b/jeecg-module-weather/src/main/java/org/jeecg/service/WeatherDataService.java @@ -34,10 +34,6 @@ public interface WeatherDataService extends IService { */ void delete(List ids); - /** - * 处理静态气象数据入库接口,比上传快 - */ - void handleStaticDataToDB(String path,Integer dataSource); /** * 关联气象数据 * @param dataType diff --git a/jeecg-module-weather/src/main/java/org/jeecg/service/WeatherTaskService.java b/jeecg-module-weather/src/main/java/org/jeecg/service/WeatherTaskService.java index e8ac6b9..1f04300 100644 --- a/jeecg-module-weather/src/main/java/org/jeecg/service/WeatherTaskService.java +++ b/jeecg-module-weather/src/main/java/org/jeecg/service/WeatherTaskService.java @@ -61,4 +61,16 @@ public interface WeatherTaskService extends IService { * @return */ List getTaskLog(Integer taskId); + + /** + * 设置任务置顶 + * @param taskId + */ + void setTaskTop(Integer taskId); + + /** + * 取消置顶 + * @param taskId + */ + void setCancelTaskTop(Integer taskId); } diff --git a/jeecg-module-weather/src/main/java/org/jeecg/service/impl/WeatherDataServiceImpl.java b/jeecg-module-weather/src/main/java/org/jeecg/service/impl/WeatherDataServiceImpl.java index 9836d0d..c336b23 100644 --- a/jeecg-module-weather/src/main/java/org/jeecg/service/impl/WeatherDataServiceImpl.java +++ b/jeecg-module-weather/src/main/java/org/jeecg/service/impl/WeatherDataServiceImpl.java @@ -358,48 +358,6 @@ public class WeatherDataServiceImpl extends ServiceImpl queryWrapper = new LambdaQueryWrapper<>(); - queryWrapper.eq(WeatherData::getMd5Value, gribFileMD5); - WeatherData queryResult = this.baseMapper.selectOne(queryWrapper); - if (Objects.nonNull(queryResult)) { - weatherLinkedDataLogService.create(dataType,file.getAbsolutePath()+"已存在"); - continue; - } - WeatherData weatherData = new WeatherData(); - weatherData.setFileName(file.getName()); - weatherData.setFileExt(file.getName().substring(file.getName().lastIndexOf(".")+1)); - weatherData.setDataSource(dataType); - weatherData.setFilePath(file.getAbsolutePath()); //获取文件数据开始日期 LocalDateTime localDateTime = Grib2TimeReader.readValidTime(file.getAbsolutePath()); - weatherData.setDataStartTime(localDateTime); - this.baseMapper.insert(weatherData); - - weatherLinkedDataLogService.create(dataType,file.getAbsolutePath()+"关联成功"); + LocalDate localDate = localDateTime.toLocalDate(); + //再次验证grib文件的有效数据时间在给定时间内,才关联保存 + if(!localDate.isBefore(startDate) && !localDate.isAfter(endDate)) { + //如果此文件存在则无需再次新增 + String gribFileMD5 = this.getGribFileMD5(file.getAbsolutePath()); + LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); + queryWrapper.eq(WeatherData::getMd5Value, gribFileMD5); + WeatherData queryResult = this.baseMapper.selectOne(queryWrapper); + if (Objects.nonNull(queryResult)) { + weatherLinkedDataLogService.create(dataType,file.getAbsolutePath()+"已存在"); + continue; + } + WeatherData weatherData = new WeatherData(); + weatherData.setFileName(file.getName()); + weatherData.setFileExt(file.getName().substring(file.getName().lastIndexOf(".")+1)); + weatherData.setDataSource(dataType); + weatherData.setFilePath(file.getAbsolutePath()); + weatherData.setDataStartTime(localDateTime); + weatherData.setMd5Value(gribFileMD5); + this.baseMapper.insert(weatherData); + weatherLinkedDataLogService.create(dataType,file.getAbsolutePath()+"关联成功"); + } }catch (Exception e){ log.error("关联{}气象数据文件出现错误,原因为:",file.getAbsolutePath(),e); weatherLinkedDataLogService.create(dataType,file.getAbsolutePath()+"关联失败,原因为:"+e.getMessage()); diff --git a/jeecg-module-weather/src/main/java/org/jeecg/service/impl/WeatherTaskServiceImpl.java b/jeecg-module-weather/src/main/java/org/jeecg/service/impl/WeatherTaskServiceImpl.java index bbcf0f5..e922bd3 100644 --- a/jeecg-module-weather/src/main/java/org/jeecg/service/impl/WeatherTaskServiceImpl.java +++ b/jeecg-module-weather/src/main/java/org/jeecg/service/impl/WeatherTaskServiceImpl.java @@ -10,6 +10,7 @@ import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import lombok.RequiredArgsConstructor; import org.apache.commons.lang3.StringUtils; import org.apache.logging.log4j.util.Strings; +import org.jeecg.common.constant.enums.TaskTopEnum; import org.jeecg.common.constant.enums.WeatherForecastDatasourceEnum; import org.jeecg.common.constant.enums.WeatherTaskStatusEnum; import org.jeecg.common.properties.SystemStorageProperties; @@ -31,6 +32,7 @@ import java.time.LocalDate; import java.time.LocalDateTime; import java.util.List; import java.util.Objects; +import java.util.UUID; /** * 天气预报预测任务管理 @@ -93,26 +95,15 @@ public class WeatherTaskServiceImpl extends ServiceImpl files = FileUtil.loopFiles(this.systemStorageProperties.getForecastFileTmpPath(), new FileFilter() { - @Override - public boolean accept(File file) { - String flag = "_"+weatherTask.getId(); - return file.getName().contains(flag); - } - }); - if (CollUtil.isNotEmpty(files)) { - for (File delFile : files) { - delFile.delete(); - } - } + if(FileUtil.exist(queryResult.getInputFile())){ + FileUtil.del(queryResult.getInputFile()); } } file.transferTo(storageFile); @@ -210,6 +197,9 @@ public class WeatherTaskServiceImpl extends ServiceImpl${jeecgboot.version} + + stas-cloud-consumer + + + org.springframework.boot + spring-boot-maven-plugin + + + \ No newline at end of file diff --git a/jeecg-server-cloud/jeecg-consumer-start/src/main/resources/application.yml b/jeecg-server-cloud/jeecg-consumer-start/src/main/resources/application.yml index f14ab38..5df0f24 100644 --- a/jeecg-server-cloud/jeecg-consumer-start/src/main/resources/application.yml +++ b/jeecg-server-cloud/jeecg-consumer-start/src/main/resources/application.yml @@ -1,5 +1,5 @@ server: - port: 8020 + port: 8011 spring: application: diff --git a/jeecg-server-cloud/pom.xml b/jeecg-server-cloud/pom.xml index 0ceb0d8..3993014 100644 --- a/jeecg-server-cloud/pom.xml +++ b/jeecg-server-cloud/pom.xml @@ -22,6 +22,7 @@ jeecg-large-screen-start jeecg-visual jeecg-weather-start + jeecg-consumer-start jeecg-event-start jeecg-sync-start jeecg-data-analyze-start diff --git a/pom.xml b/pom.xml index b5eb6cf..33e3cd8 100644 --- a/pom.xml +++ b/pom.xml @@ -93,7 +93,6 @@ jeecg-module-transport jeecg-module-monitor-info-database jeecg-model-consumer - jeecg-server-cloud/jeecg-consumer-start