From e477e90099822c3f411208fae5e9583309d088ea Mon Sep 17 00:00:00 2001 From: panbaolin Date: Mon, 20 Jul 2026 17:49:30 +0800 Subject: [PATCH] =?UTF-8?q?fix:1.=E6=B7=BB=E5=8A=A0gfs=E6=95=B0=E6=8D=AE?= =?UTF-8?q?=E4=B8=8B=E8=BD=BD=E5=8A=9F=E8=83=BD2.=E6=B7=BB=E5=8A=A0gfs?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=E5=92=8Ct1h=E6=95=B0=E6=8D=AE=E5=AE=9A?= =?UTF-8?q?=E6=97=B6=E6=B8=85=E9=99=A4=E5=8A=9F=E8=83=BD3.=E4=BC=98?= =?UTF-8?q?=E5=8C=96=E5=85=B6=E4=BB=96=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../constant/WeatherPrefixConstants.java | 10 - .../constant/WeatherSuffixConstants.java | 7 - .../constant/enums/WeatherDataSourceEnum.java | 6 +- .../enums/WeatherVariableNameEnum.java | 7 +- .../properties/SystemStorageProperties.java | 5 + .../TransportSimulationProperties.java | 16 + .../jeecg/common/util/Grib2TimeReader.java | 2 +- .../jeecg/common/util}/JSchRemoteRunner.java | 13 +- .../modules/base/entity/WeatherData.java | 2 +- .../base/entity/WeatherDownGFSDataLog.java | 43 +++ .../base/entity/WeatherDownT1hDataLog.java | 42 +++ .../mapper/WeatherDownGFSDataLogMapper.java | 8 + .../mapper/WeatherDownT1hDataLogMapper.java | 7 + .../base/mapper/xml/TransportTaskMapper.xml | 3 + .../rebuild/task/SourceRebuildTaskExec.java | 2 +- .../china/AbstractTaskMsgHandler.java | 1 - .../flexparttask/AbstractTaskExec.java | 29 +- .../flexparttask/BackwardTaskExec.java | 5 +- .../flexparttask/ForwardTaskExec.java | 4 +- .../weather/task/WeatherForecastTaskExec.java | 2 +- .../controller/StasSyncLogController.java | 11 +- .../syncLog/service/IStasSyncLogService.java | 5 +- .../service/impl/StasSyncLogServiceImpl.java | 10 +- .../TransportResultDataController.java | 7 + .../service/TransportResultDataService.java | 9 + .../impl/TransportResultDataServiceImpl.java | 65 +++++ .../impl/TransportTaskServiceImpl.java | 3 +- .../controller/WeatherDataController.java | 4 +- .../main/java/org/jeecg/job/CleanDataJob.java | 209 +++++++++++++ .../java/org/jeecg/job/DownloadGFSJob.java | 274 ++++++++++++++++++ .../service/impl/WeatherDataServiceImpl.java | 8 + 31 files changed, 767 insertions(+), 52 deletions(-) delete mode 100644 jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/WeatherPrefixConstants.java delete mode 100644 jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/WeatherSuffixConstants.java rename {jeecg-model-consumer/src/main/java/org/jeecg/jsch => jeecg-boot-base-core/src/main/java/org/jeecg/common/util}/JSchRemoteRunner.java (94%) create mode 100644 jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherDownGFSDataLog.java create mode 100644 jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherDownT1hDataLog.java create mode 100644 jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/WeatherDownGFSDataLogMapper.java create mode 100644 jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/WeatherDownT1hDataLogMapper.java create mode 100644 jeecg-module-weather/src/main/java/org/jeecg/job/CleanDataJob.java create mode 100644 jeecg-module-weather/src/main/java/org/jeecg/job/DownloadGFSJob.java diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/WeatherPrefixConstants.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/WeatherPrefixConstants.java deleted file mode 100644 index b805b22..0000000 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/WeatherPrefixConstants.java +++ /dev/null @@ -1,10 +0,0 @@ -package org.jeecg.common.constant; - -public class WeatherPrefixConstants { - - public static final String PANGU_PREFIX = "panguweather_"; - public static final String CRA40_PREFIX = "CRA40_"; - public static final String NCEP_PREFIX = "cdas1"; - public static final String T1H_PREFIX = "T1H_"; - public static final String FNL_PREFIX = "fnl_"; -} diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/WeatherSuffixConstants.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/WeatherSuffixConstants.java deleted file mode 100644 index 13ad22f..0000000 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/WeatherSuffixConstants.java +++ /dev/null @@ -1,7 +0,0 @@ -package org.jeecg.common.constant; - -public class WeatherSuffixConstants { - - public static final String CRA40_SUFFIX = "_GLB_0P25_HOUR_V1_0_0"; - public static final String NCEP_SUFFIX = "pgrbh"; -} diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/WeatherDataSourceEnum.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/WeatherDataSourceEnum.java index 96f3233..dfb5db2 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/WeatherDataSourceEnum.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/WeatherDataSourceEnum.java @@ -31,7 +31,11 @@ public enum WeatherDataSourceEnum { /** * T1H */ - T1H(6,"T1H"); + T1H(6,"T1H"), + /** + * GFS + */ + GFS(7,"GFS"); private Integer key; private String value; diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/WeatherVariableNameEnum.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/WeatherVariableNameEnum.java index 08dbbaf..79bb873 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/WeatherVariableNameEnum.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/common/constant/enums/WeatherVariableNameEnum.java @@ -30,7 +30,12 @@ public enum WeatherVariableNameEnum { T1H_P(WeatherDataSourceEnum.T1H.getKey(), 1, "Pressure_height_above_ground"), T1H_H(WeatherDataSourceEnum.T1H.getKey(), 2, "Relative_humidity_height_above_ground"), T1H_U(WeatherDataSourceEnum.T1H.getKey(), 3, "u-component_of_wind_height_above_ground"), - T1H_V(WeatherDataSourceEnum.T1H.getKey(), 4, "v-component_of_wind_height_above_ground"); + T1H_V(WeatherDataSourceEnum.T1H.getKey(), 4, "v-component_of_wind_height_above_ground"), + GFS_T(WeatherDataSourceEnum.GFS.getKey(), 0, "Temperature_height_above_ground"), + GFS_P(WeatherDataSourceEnum.GFS.getKey(), 1, "Pressure_height_above_ground"), + GFS_H(WeatherDataSourceEnum.GFS.getKey(), 2, "Relative_humidity_height_above_ground"), + GFS_U(WeatherDataSourceEnum.GFS.getKey(), 3, "u-component_of_wind_height_above_ground"), + GFS_V(WeatherDataSourceEnum.GFS.getKey(), 4, "v-component_of_wind_height_above_ground"); private Integer type; diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/common/properties/SystemStorageProperties.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/properties/SystemStorageProperties.java index 940b33f..a47d7aa 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/common/properties/SystemStorageProperties.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/common/properties/SystemStorageProperties.java @@ -34,6 +34,11 @@ public class SystemStorageProperties { */ private String t1hDataPath; + /** + * gfs数据存储路径 + */ + private String gfsDataPath; + /** * fnl数据存储路径 */ diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/common/properties/TransportSimulationProperties.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/properties/TransportSimulationProperties.java index 2dfacf0..e37260d 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/common/properties/TransportSimulationProperties.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/common/properties/TransportSimulationProperties.java @@ -53,6 +53,10 @@ public class TransportSimulationProperties { * 参数配置文件路径(和T1H气象数据有关系) */ private String t1hParamConfigPath; + /** + * 参数配置文件路径(和GFS气象数据有关系) + */ + private String gfsParamConfigPath; /** * 台站数据配置文件路径(和fnl气象数据有关系) @@ -78,6 +82,10 @@ public class TransportSimulationProperties { * 台站数据配置文件路径(和T1H气象数据有关系) */ private String t1hStationsConfigPath; + /** + * 台站数据配置文件路径(和GFS气象数据有关系) + */ + private String gfsStationsConfigPath; /** * 正演脚本路径(和fnl气象数据有关系) @@ -103,6 +111,10 @@ public class TransportSimulationProperties { * 正演脚本路径(和T1H气象数据有关系) */ private String t1hForwardScriptPath; + /** + * 正演脚本路径(和GFS气象数据有关系) + */ + private String gfsForwardScriptPath; /** * 反演脚本路径(和fnl气象数据有关系) @@ -128,6 +140,10 @@ public class TransportSimulationProperties { * 反演脚本路径(和T1H气象数据有关系) */ private String t1hBackwardScriptPath; + /** + * 反演脚本路径(和GFS气象数据有关系) + */ + private String gfsBackwardScriptPath; /** * nc转换gif配置 diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/common/util/Grib2TimeReader.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/util/Grib2TimeReader.java index dbd9bc0..0faa785 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/common/util/Grib2TimeReader.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/common/util/Grib2TimeReader.java @@ -70,7 +70,7 @@ public class Grib2TimeReader { } public static void main(String[] args) throws Exception { - LocalDateTime validTime = readValidTime("C:\\Users\\cnndc\\Desktop\\pangu_20260605_12_00.grib2"); + LocalDateTime validTime = readValidTime("C:\\Users\\cnndc\\Desktop\\gfs.t00z.pgrb2.0p25.f004"); String result = validTime.format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")); System.out.println(result); diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/jsch/JSchRemoteRunner.java b/jeecg-boot-base-core/src/main/java/org/jeecg/common/util/JSchRemoteRunner.java similarity index 94% rename from jeecg-model-consumer/src/main/java/org/jeecg/jsch/JSchRemoteRunner.java rename to jeecg-boot-base-core/src/main/java/org/jeecg/common/util/JSchRemoteRunner.java index ddaed43..72c6a82 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/jsch/JSchRemoteRunner.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/common/util/JSchRemoteRunner.java @@ -1,4 +1,4 @@ -package org.jeecg.jsch; +package org.jeecg.common.util; import cn.hutool.core.collection.CollUtil; import cn.hutool.core.util.StrUtil; @@ -6,7 +6,8 @@ import com.jcraft.jsch.*; import lombok.Getter; import lombok.extern.slf4j.Slf4j; -import java.io.*; +import java.io.ByteArrayInputStream; +import java.io.File; import java.nio.charset.StandardCharsets; import java.util.*; @@ -193,7 +194,9 @@ public class JSchRemoteRunner { List matchedFiles = new ArrayList<>(); for (ChannelSftp.LsEntry entry : fileList) { String fileName = entry.getFilename(); - if (fileName.endsWith(matchingSuffix)) { + if (StrUtil.isNotBlank(matchingSuffix) && fileName.endsWith(matchingSuffix)) { + matchedFiles.add(new File(path+File.separator+fileName)); + }else if(StrUtil.isBlank(matchingSuffix)){ matchedFiles.add(new File(path+File.separator+fileName)); } } @@ -223,7 +226,9 @@ public class JSchRemoteRunner { Vector fileList = sftp.ls(path); for (ChannelSftp.LsEntry entry : fileList) { String fileName = entry.getFilename(); - if (fileName.endsWith(matchingSuffix)) { + if (StrUtil.isNotBlank(matchingSuffix) && fileName.endsWith(matchingSuffix)) { + matchedFiles.add(path+File.separator+fileName); + }else if(StrUtil.isBlank(matchingSuffix)){ matchedFiles.add(path+File.separator+fileName); } } diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherData.java b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherData.java index b1cd39a..ae18cd4 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherData.java +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherData.java @@ -47,7 +47,7 @@ public class WeatherData implements Serializable { private LocalDateTime dataStartTime; /** - * 数据来源(1-盘古模型,2-graphcast,3-cra40,4-ncep,5-fnl,6-t1h) + * 数据来源(1-盘古模型,2-graphcast,3-cra40,4-ncep,5-fnl,6-t1h,7-gfs) */ @TableField(value = "data_source") private Integer dataSource; diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherDownGFSDataLog.java b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherDownGFSDataLog.java new file mode 100644 index 0000000..d7955c7 --- /dev/null +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherDownGFSDataLog.java @@ -0,0 +1,43 @@ +package org.jeecg.modules.base.entity; + +import com.baomidou.mybatisplus.annotation.IdType; +import com.baomidou.mybatisplus.annotation.TableField; +import com.baomidou.mybatisplus.annotation.TableId; +import com.baomidou.mybatisplus.annotation.TableName; +import com.fasterxml.jackson.annotation.JsonFormat; +import lombok.Data; + +import java.time.LocalDateTime; + +/** + * 记录下载GFS气象数据日志数据表 + */ +@Data +@TableName("stas_gfs_download_log") +public class WeatherDownGFSDataLog { + + /** + * ID + */ + @TableId(type = IdType.AUTO) + private Integer id; + + /** + * 批次时间 + */ + @TableField(value = "batch_time") + private String batchTime; + + /** + * 创建时间 + */ + @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss") + @TableField(value = "create_time") + private LocalDateTime createTime; + + /** + * 任务运行过程日志 + */ + @TableField(value = "log_content") + private String logContent; +} diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherDownT1hDataLog.java b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherDownT1hDataLog.java new file mode 100644 index 0000000..3436ec9 --- /dev/null +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/entity/WeatherDownT1hDataLog.java @@ -0,0 +1,42 @@ +package org.jeecg.modules.base.entity; + +import com.baomidou.mybatisplus.annotation.IdType; +import com.baomidou.mybatisplus.annotation.TableField; +import com.baomidou.mybatisplus.annotation.TableId; +import com.baomidou.mybatisplus.annotation.TableName; +import com.fasterxml.jackson.annotation.JsonFormat; +import lombok.Data; +import java.time.LocalDateTime; + +/** + * 记录下载T1h气象数据日志数据表 + */ +@Data +@TableName("stas_t1h_download_log") +public class WeatherDownT1hDataLog { + + /** + * ID + */ + @TableId(type = IdType.AUTO) + private Integer id; + + /** + * 批次时间 + */ + @TableField(value = "batch_time") + private String batchTime; + + /** + * 创建时间 + */ + @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss") + @TableField(value = "create_time") + private LocalDateTime createTime; + + /** + * 任务运行过程日志 + */ + @TableField(value = "log_content") + private String logContent; +} diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/WeatherDownGFSDataLogMapper.java b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/WeatherDownGFSDataLogMapper.java new file mode 100644 index 0000000..2e6200c --- /dev/null +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/WeatherDownGFSDataLogMapper.java @@ -0,0 +1,8 @@ +package org.jeecg.modules.base.mapper; + +import com.baomidou.mybatisplus.core.mapper.BaseMapper; +import org.jeecg.modules.base.entity.WeatherDownGFSDataLog; + +public interface WeatherDownGFSDataLogMapper extends BaseMapper { + +} diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/WeatherDownT1hDataLogMapper.java b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/WeatherDownT1hDataLogMapper.java new file mode 100644 index 0000000..1003394 --- /dev/null +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/WeatherDownT1hDataLogMapper.java @@ -0,0 +1,7 @@ +package org.jeecg.modules.base.mapper; + +import com.baomidou.mybatisplus.core.mapper.BaseMapper; +import org.jeecg.modules.base.entity.WeatherDownT1hDataLog; + +public interface WeatherDownT1hDataLogMapper extends BaseMapper { +} diff --git a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/xml/TransportTaskMapper.xml b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/xml/TransportTaskMapper.xml index c4be159..6a72cee 100644 --- a/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/xml/TransportTaskMapper.xml +++ b/jeecg-boot-base-core/src/main/java/org/jeecg/modules/base/mapper/xml/TransportTaskMapper.xml @@ -15,6 +15,9 @@ t.top_task as topTask from stas_transport_task t + + t.task_name like '%' || ${taskName} || '%' + t.task_mode = #{taskMode} diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/task/SourceRebuildTaskExec.java b/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/task/SourceRebuildTaskExec.java index be57ce3..0416c76 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/task/SourceRebuildTaskExec.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/rebuild/task/SourceRebuildTaskExec.java @@ -9,7 +9,7 @@ import org.jeecg.common.constant.enums.SourceRebuildReleaseSourceEnum; import org.jeecg.common.constant.enums.SourceRebuildTaskStatusEnum; import org.jeecg.common.properties.ServerProperties; import org.jeecg.common.properties.SourceRebuildProperties; -import org.jeecg.jsch.JSchRemoteRunner; +import org.jeecg.common.util.JSchRemoteRunner; import org.jeecg.modules.base.entity.SourceRebuildMonitoringData; import org.jeecg.modules.base.entity.SourceRebuildTask; import org.jeecg.modules.base.entity.SourceRebuildTaskLog; diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/transport/consumer/china/AbstractTaskMsgHandler.java b/jeecg-model-consumer/src/main/java/org/jeecg/transport/consumer/china/AbstractTaskMsgHandler.java index c525fcf..9f248b3 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/transport/consumer/china/AbstractTaskMsgHandler.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/transport/consumer/china/AbstractTaskMsgHandler.java @@ -90,7 +90,6 @@ public abstract class AbstractTaskMsgHandler extends AbstractChain{ * @param transportTask */ protected void runTask(TransportTask transportTask,String ip){ - log.info("收到任务:"+transportTask.getTaskName()); boolean flag = false; if (TransportTaskModeEnum.BACK_FORWARD.getKey().equals(transportTask.getTaskMode())){ AbstractTaskExec taskExec = new BackwardTaskExec(); diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/AbstractTaskExec.java b/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/AbstractTaskExec.java index dbd52fb..f4ca828 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/AbstractTaskExec.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/transport/flexparttask/AbstractTaskExec.java @@ -27,6 +27,7 @@ import org.jeecg.transport.service.StationsModValService; import org.jeecg.transport.service.TransportTaskService; import java.io.*; import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; import java.time.temporal.ChronoUnit; import java.util.ArrayList; import java.util.List; @@ -170,7 +171,6 @@ public abstract class AbstractTaskExec extends Thread{ command.append(File.separator); command.append(ncFile.getParentFile().getName()); command.append(".gif"); - log.info(command.toString()); this.execGenGIFCommand(command.toString(),session); } }catch (Exception e){ @@ -257,7 +257,26 @@ public abstract class AbstractTaskExec extends Thread{ }else if(WeatherDataSourceEnum.FNL.getKey().equals(dataSource)){ path.append(systemStorageProperties.getFnlDataPath()); }else if(WeatherDataSourceEnum.T1H.getKey().equals(dataSource)){ - path.append(systemStorageProperties.getT1hDataPath()); + //T1H是中科天机的预报数据,正演只能按批次来,最大15天,每批次存储当前天往后15天的预报数据 + if(TransportTaskModeEnum.FORWARD.getKey().equals(transportTask.getTaskMode())){ + DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyyMMdd"); + String format = formatter.format(transportTask.getStartTime()); + path.append(File.separator); + path.append(format); + }else { + path.append(systemStorageProperties.getT1hDataPath()); + } + }else if(WeatherDataSourceEnum.GFS.getKey().equals(dataSource)){ + //GFS是美国的预报数据,正演只能按批次来,最大16天,每批次存储当前天往后16天的预报数据 + //这里补一下数据所在批次的目录 + if(TransportTaskModeEnum.FORWARD.getKey().equals(transportTask.getTaskMode())){ + DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyyMMdd_HH"); + String format = formatter.format(transportTask.getStartTime()); + path.append(File.separator); + path.append("GFS_"+format); + }else { + path.append(systemStorageProperties.getGfsDataPath()); + } } return path.toString(); } @@ -280,6 +299,8 @@ public abstract class AbstractTaskExec extends Thread{ return simulationProperties.getFnlParamConfigPath(); }else if(WeatherDataSourceEnum.T1H.getKey().equals(dataSource)){ return simulationProperties.getT1hParamConfigPath(); + }else if(WeatherDataSourceEnum.GFS.getKey().equals(dataSource)){ + return simulationProperties.getGfsParamConfigPath(); } return Strings.EMPTY; } @@ -294,7 +315,7 @@ public abstract class AbstractTaskExec extends Thread{ if(WeatherDataSourceEnum.PANGU.getKey().equals(dataSource)){ path.append(simulationProperties.getPanguStationsConfigPath()); }else if(WeatherDataSourceEnum.GRAPHCAST.getKey().equals(dataSource)){ - path.append(simulationProperties.getCra40StationsConfigPath()); + path.append(simulationProperties.getGraphcastStationsConfigPath()); }else if(WeatherDataSourceEnum.CRA40.getKey().equals(dataSource)){ path.append(simulationProperties.getCra40StationsConfigPath()); }else if(WeatherDataSourceEnum.NCEP.getKey().equals(dataSource)){ @@ -303,6 +324,8 @@ public abstract class AbstractTaskExec extends Thread{ path.append(simulationProperties.getFnlStationsConfigPath()); }else if(WeatherDataSourceEnum.T1H.getKey().equals(dataSource)){ path.append(simulationProperties.getT1hStationsConfigPath()); + }else if(WeatherDataSourceEnum.GFS.getKey().equals(dataSource)){ + path.append(simulationProperties.getGfsStationsConfigPath()); } return path.toString(); } 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 dd74dc5..ca93d1c 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 @@ -8,9 +8,8 @@ import org.apache.logging.log4j.util.Strings; import org.jeecg.common.constant.enums.FlexpartSpeciesType; import org.jeecg.common.constant.enums.TransportTaskStatusEnum; import org.jeecg.common.constant.enums.WeatherDataSourceEnum; +import org.jeecg.common.util.JSchRemoteRunner; import org.jeecg.modules.base.entity.*; -import org.jeecg.jsch.JSchRemoteRunner; - import java.io.File; import java.math.BigDecimal; import java.math.RoundingMode; @@ -207,6 +206,8 @@ public class BackwardTaskExec extends AbstractTaskExec { return simulationProperties.getFnlBackwardScriptPath(); }else if(WeatherDataSourceEnum.T1H.getKey().equals(dataSource)){ return simulationProperties.getT1hBackwardScriptPath(); + }else if(WeatherDataSourceEnum.GFS.getKey().equals(dataSource)){ + return simulationProperties.getGfsBackwardScriptPath(); } return Strings.EMPTY; } 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 fc3859a..54b1f8d 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 @@ -7,8 +7,8 @@ import org.apache.commons.lang3.time.StopWatch; import org.apache.logging.log4j.util.Strings; import org.jeecg.common.constant.enums.TransportTaskStatusEnum; import org.jeecg.common.constant.enums.WeatherDataSourceEnum; +import org.jeecg.common.util.JSchRemoteRunner; import org.jeecg.modules.base.entity.*; -import org.jeecg.jsch.JSchRemoteRunner; import java.io.File; import java.math.BigDecimal; import java.math.RoundingMode; @@ -269,6 +269,8 @@ public class ForwardTaskExec extends AbstractTaskExec { return simulationProperties.getFnlForwardScriptPath(); }else if(WeatherDataSourceEnum.T1H.getKey().equals(dataSource)){ return simulationProperties.getT1hForwardScriptPath(); + }else if(WeatherDataSourceEnum.GFS.getKey().equals(dataSource)){ + return simulationProperties.getGfsForwardScriptPath(); } return Strings.EMPTY; } diff --git a/jeecg-model-consumer/src/main/java/org/jeecg/weather/task/WeatherForecastTaskExec.java b/jeecg-model-consumer/src/main/java/org/jeecg/weather/task/WeatherForecastTaskExec.java index 6c75b1d..54252d7 100644 --- a/jeecg-model-consumer/src/main/java/org/jeecg/weather/task/WeatherForecastTaskExec.java +++ b/jeecg-model-consumer/src/main/java/org/jeecg/weather/task/WeatherForecastTaskExec.java @@ -12,7 +12,7 @@ import org.jeecg.common.constant.enums.WeatherFileSuffixEnum; import org.jeecg.common.constant.enums.WeatherForecastDatasourceEnum; import org.jeecg.common.constant.enums.WeatherTaskStatusEnum; import org.jeecg.common.util.Grib2TimeReader; -import org.jeecg.jsch.JSchRemoteRunner; +import org.jeecg.common.util.JSchRemoteRunner; import org.jeecg.modules.base.entity.WeatherData; import org.springframework.http.MediaType; import org.springframework.web.reactive.function.client.WebClient; diff --git a/jeecg-module-sync/src/main/java/org/jeecg/syncLog/controller/StasSyncLogController.java b/jeecg-module-sync/src/main/java/org/jeecg/syncLog/controller/StasSyncLogController.java index 40ee3a4..ce573f8 100644 --- a/jeecg-module-sync/src/main/java/org/jeecg/syncLog/controller/StasSyncLogController.java +++ b/jeecg-module-sync/src/main/java/org/jeecg/syncLog/controller/StasSyncLogController.java @@ -5,11 +5,12 @@ import io.swagger.v3.oas.annotations.tags.Tag; import org.jeecg.common.api.vo.Result; import org.jeecg.modules.base.entity.StasSyncLog; import org.jeecg.syncLog.service.IStasSyncLogService; -import com.baomidou.mybatisplus.core.metadata.IPage; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.*; +import java.util.List; + /** * 同步日志信息 */ @@ -26,15 +27,13 @@ public class StasSyncLogController{ * 分页列表查询 * * @param recordId - * @param pageNum - * @param pageSize * @return */ @Operation(summary = "同步日志信息-分页列表查询") @GetMapping(value = "/list") - public Result> queryPageList(String recordId, Integer pageNum, Integer pageSize) { - IPage pageList = stasSyncLogService.queryPageList(recordId, pageNum,pageSize); - return Result.OK(pageList); + public Result queryList(String recordId) { + List list = stasSyncLogService.queryList(recordId); + return Result.OK(list); } } diff --git a/jeecg-module-sync/src/main/java/org/jeecg/syncLog/service/IStasSyncLogService.java b/jeecg-module-sync/src/main/java/org/jeecg/syncLog/service/IStasSyncLogService.java index 6a6733e..af7a0f7 100644 --- a/jeecg-module-sync/src/main/java/org/jeecg/syncLog/service/IStasSyncLogService.java +++ b/jeecg-module-sync/src/main/java/org/jeecg/syncLog/service/IStasSyncLogService.java @@ -5,6 +5,7 @@ import org.jeecg.modules.base.entity.StasSyncLog; import com.baomidou.mybatisplus.extension.service.IService; import java.time.LocalDateTime; +import java.util.List; /** * 同步日志信息 @@ -14,9 +15,7 @@ public interface IStasSyncLogService extends IService { /** * 分页查询同步日志 * @param recordId - * @param pageNum - * @param pageSize * @return */ - IPage queryPageList(String recordId, Integer pageNum, Integer pageSize); + List queryList(String recordId); } diff --git a/jeecg-module-sync/src/main/java/org/jeecg/syncLog/service/impl/StasSyncLogServiceImpl.java b/jeecg-module-sync/src/main/java/org/jeecg/syncLog/service/impl/StasSyncLogServiceImpl.java index fc976f3..20494b7 100644 --- a/jeecg-module-sync/src/main/java/org/jeecg/syncLog/service/impl/StasSyncLogServiceImpl.java +++ b/jeecg-module-sync/src/main/java/org/jeecg/syncLog/service/impl/StasSyncLogServiceImpl.java @@ -9,6 +9,8 @@ import org.jeecg.syncLog.service.IStasSyncLogService; import org.springframework.stereotype.Service; import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; +import java.util.List; + /** * 同步日志信息 */ @@ -19,17 +21,13 @@ public class StasSyncLogServiceImpl extends ServiceImpl queryPageList(String recordId, Integer pageNum, Integer pageSize) { + public List queryList(String recordId) { LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); queryWrapper.eq(StasSyncLog::getRecordId, recordId); queryWrapper.orderByAsc(StasSyncLog::getStartTime); - - IPage page = new Page<>(pageNum, pageSize); - return this.baseMapper.selectPage(page, queryWrapper); + return this.baseMapper.selectList(queryWrapper); } } diff --git a/jeecg-module-transport/src/main/java/org/jeecg/controller/TransportResultDataController.java b/jeecg-module-transport/src/main/java/org/jeecg/controller/TransportResultDataController.java index 6e1fbcc..39b8096 100644 --- a/jeecg-module-transport/src/main/java/org/jeecg/controller/TransportResultDataController.java +++ b/jeecg-module-transport/src/main/java/org/jeecg/controller/TransportResultDataController.java @@ -115,4 +115,11 @@ public class TransportResultDataController { public Result getTaskStations(@NotNull(message = "任务id不能为空") Integer taskId) { return Result.OK(transportResultDataService.getTaskStations(taskId)); } + + @AutoLog(value = "查询任务所属GIF文件地址") + @Operation(summary = "查询任务所属GIF文件地址") + @GetMapping("getTaskGifAddr") + public Result getTaskGifAddr(@NotNull(message = "任务id不能为空") Integer taskId,Integer speciesId) { + return Result.OK(transportResultDataService.getTaskGifAddr(taskId,speciesId)); + } } diff --git a/jeecg-module-transport/src/main/java/org/jeecg/service/TransportResultDataService.java b/jeecg-module-transport/src/main/java/org/jeecg/service/TransportResultDataService.java index 90e5c1a..82b12fa 100644 --- a/jeecg-module-transport/src/main/java/org/jeecg/service/TransportResultDataService.java +++ b/jeecg-module-transport/src/main/java/org/jeecg/service/TransportResultDataService.java @@ -1,5 +1,6 @@ package org.jeecg.service; +import jakarta.validation.constraints.NotNull; import org.jeecg.vo.ContributionAnalysisVO; import org.jeecg.vo.QueryDiffusionVO; import org.jeecg.vo.TaskStationsVO; @@ -108,4 +109,12 @@ public interface TransportResultDataService { * @return */ Map>> getConformityAnalysisScatterPlot(Integer taskId,String stationIds,Integer speciesId,String nuclideName); + + /** + * 获取gif地址 + * @param taskId + * @param speciesId + * @return + */ + String getTaskGifAddr(Integer taskId,Integer speciesId); } diff --git a/jeecg-module-transport/src/main/java/org/jeecg/service/impl/TransportResultDataServiceImpl.java b/jeecg-module-transport/src/main/java/org/jeecg/service/impl/TransportResultDataServiceImpl.java index 8133182..4da1298 100644 --- a/jeecg-module-transport/src/main/java/org/jeecg/service/impl/TransportResultDataServiceImpl.java +++ b/jeecg-module-transport/src/main/java/org/jeecg/service/impl/TransportResultDataServiceImpl.java @@ -4,6 +4,7 @@ import cn.hutool.core.collection.CollUtil; import cn.hutool.core.date.DateUtil; import cn.hutool.core.date.LocalDateTimeUtil; import cn.hutool.core.io.FileUtil; +import cn.hutool.core.util.StrUtil; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -674,6 +675,28 @@ public class TransportResultDataServiceImpl implements TransportResultDataServic return resultMap; } + /** + * 获取gif地址 + * @param taskId + * @param speciesId + * @return + */ + @Override + public String getTaskGifAddr(Integer taskId,Integer speciesId) { + TransportTask transportTask = this.transportTaskMapper.selectById(taskId); + //nginx代理父路径就是/gif + String parentPath = "/gif"; + if(TransportTaskModeEnum.FORWARD.getKey().equals(transportTask.getTaskMode())){ + return this.getForwardTaskGIFPath(transportTask,speciesId,parentPath); + }else if(TransportTaskModeEnum.BACK_FORWARD.getKey().equals(transportTask.getTaskMode())){ + LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); + queryWrapper.eq(TransportTaskBackwardChild::getTaskId,taskId); + List taskBackwardChildren = this.backwardChildMapper.selectList(queryWrapper); + return this.getBackForwardTaskGIFPath(transportTask,taskBackwardChildren.get(0).getStationCode(),parentPath); + } + return StrUtil.EMPTY; + } + /** * 获取正演NC结果文件路径 * @param transportTask @@ -694,6 +717,27 @@ public class TransportResultDataServiceImpl implements TransportResultDataServic return path.toString(); } + /** + * 获取正演GIF文件路径 + * @param transportTask + * @param speciesId + * @return + */ + private String getForwardTaskGIFPath(TransportTask transportTask,Integer speciesId,String parentPath){ + //拼接GIF文件路径 + StringBuilder path = new StringBuilder(); + path.append(parentPath); + path.append(File.separator); + path.append(transportTask.getTaskName()); + path.append(File.separator); + path.append(FORWARD); + path.append(File.separator); + path.append(speciesId.toString()); + path.append(File.separator); + path.append(speciesId+".gif"); + return path.toString(); + } + /** * 获取反演NC结果文件路径 * @param transportTask @@ -714,6 +758,27 @@ public class TransportResultDataServiceImpl implements TransportResultDataServic return path.toString(); } + /** + * 获取反演GIF文件路径 + * @param transportTask + * @param stationCode + * @return + */ + private String getBackForwardTaskGIFPath(TransportTask transportTask,String stationCode,String parentPath){ + //拼接gif文件路径 + StringBuilder path = new StringBuilder(); + path.append(parentPath); + path.append(File.separator); + path.append(transportTask.getTaskName()); + path.append(File.separator); + path.append(BACK_FORWARD); + path.append(File.separator); + path.append(stationCode); + path.append(File.separator); + path.append(stationCode+".gif"); + return path.toString(); + } + /** * 在步长为 0.25° 的递增经纬度数组中,查找指定位置的索引 * - lons[0] 是起始经度 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 2584afd..f3d2f86 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 @@ -1107,7 +1107,8 @@ public class TransportTaskServiceImpl extends ServiceImpl getTimeBatch() { + public Result getTimeBatch(Integer dataType) { List timeBatchList = weatherDataService .list(new LambdaQueryWrapper() .select(WeatherData::getTimeBatch) // 只查询需要的字段 - .eq(WeatherData::getDataSource, WeatherDataSourceEnum.T1H.getKey())) + .eq(WeatherData::getDataSource,dataType)) .stream() .map(WeatherData::getTimeBatch) .distinct() diff --git a/jeecg-module-weather/src/main/java/org/jeecg/job/CleanDataJob.java b/jeecg-module-weather/src/main/java/org/jeecg/job/CleanDataJob.java new file mode 100644 index 0000000..6407927 --- /dev/null +++ b/jeecg-module-weather/src/main/java/org/jeecg/job/CleanDataJob.java @@ -0,0 +1,209 @@ +package org.jeecg.job; + +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.baomidou.mybatisplus.core.toolkit.StringPool; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.jeecg.modules.base.entity.WeatherDownGFSDataLog; +import org.jeecg.modules.base.entity.WeatherDownT1hDataLog; +import org.jeecg.modules.base.mapper.WeatherDownGFSDataLogMapper; +import org.jeecg.modules.base.mapper.WeatherDownT1hDataLogMapper; +import org.jetbrains.annotations.NotNull; +import org.jetbrains.annotations.Nullable; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +import java.io.File; +import java.io.IOException; +import java.nio.file.*; +import java.nio.file.attribute.BasicFileAttributes; +import java.time.LocalDate; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Stream; + +/** + * 每天定时清除T1h和GFS历史数据 + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class CleanDataJob { + + @Value("${gfs-download.cleanSwitch}") + private boolean gfsCleanSwitch; + @Value("${gfs-download.retentionPeriod}") + private Integer gfsDataRetentionPeriod; + @Value("${gfs-download.gfsDataPath}") + private String gfsDataPath; + @Value("${t1h-download.cleanSwitch}") + private boolean t1hCleanSwitch; + @Value("${t1h-download.retentionPeriod}") + private Integer t1hDataRetentionPeriod; + @Value("${t1h-download.t1hDataPath}") + private String t1hDataPath; + private final WeatherDownGFSDataLogMapper downGFSDataLogMapper; + private final WeatherDownT1hDataLogMapper downT1hDataLogMapper; + + private boolean t1hIsFirstRun = true; + private boolean gfsIsFirstRun = true; + private final static String HOURS = "00"; + private final static String DIRPREFIX = "GFS_"; + + /** + * 每天23点清除T1H和GFS历史预报数据 + */ + @Scheduled(cron = "0 0 23 * * ?") + public void cleanDirectory(){ + if(t1hCleanSwitch){ + this.cleanT1hData(); + if (t1hIsFirstRun){ + t1hIsFirstRun = false; + } + } + if(gfsCleanSwitch){ + this.cleanGFSData(); + if (gfsIsFirstRun){ + gfsIsFirstRun = false; + } + } + } + + /** + * 清除t1h历史数据 + */ + private void cleanT1hData(){ + DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyyMMdd"); + List targetDirs = new ArrayList<>(); + if(t1hIsFirstRun){ + Path t1hDataDir = Paths.get(this.t1hDataPath); + try (Stream stream = Files.walk(t1hDataDir)){ + List dirs = stream.filter(Files::isDirectory) + .filter(path -> !path.equals(t1hDataDir)) + .toList(); + LocalDate retentionDate = LocalDate.now().minusDays(t1hDataRetentionPeriod); + for (Path dir : dirs) { + String name = dir.toFile().getName(); + LocalDate dirDate = LocalDate.parse(name, formatter); + if (!dirDate.isAfter(retentionDate)) { + targetDirs.add(dir); + } + } + } catch (IOException e) { + throw new RuntimeException(e); + } + }else { + LocalDate localDate = LocalDate.now().minusDays(t1hDataRetentionPeriod); + String batchTime = formatter.format(localDate); + String targetDir = this.t1hDataPath+File.separator+batchTime; + targetDirs.add(Paths.get(targetDir)); + } + try { + for (Path targetDirPath : targetDirs){ + if(!Files.exists(targetDirPath)){ + break; + } + Files.walkFileTree(targetDirPath,new SimpleFileVisitor<>() { + @NotNull + @Override + public FileVisitResult visitFile(@NotNull Path file, @NotNull BasicFileAttributes attrs) throws IOException { + Files.delete(file); + return FileVisitResult.CONTINUE; + } + @NotNull + @Override + public FileVisitResult postVisitDirectory(@NotNull Path dir, @Nullable IOException exc) throws IOException { + Files.delete(dir); + return FileVisitResult.CONTINUE; + } + }); + } + this.cleanT1hDataLog(); + log.info("完成{}目录清除超出保留期限气象数据",t1hDataPath); + } catch (IOException e) { + log.error("{}目录清除超出保留期限气象数据出现异常",t1hDataPath,e); + } + } + + /** + * 清除gfs历史数据 + */ + private void cleanGFSData(){ + DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyyMMdd"); + List targetDirs = new ArrayList<>(); + if(gfsIsFirstRun){ + Path gfsDataDir = Paths.get(this.gfsDataPath); + try (Stream stream = Files.walk(gfsDataDir)){ + List dirs = stream.filter(Files::isDirectory) + .filter(path -> !path.equals(gfsDataDir)) + .toList(); + LocalDate retentionDate = LocalDate.now().minusDays(gfsDataRetentionPeriod); + for (Path dir : dirs) { + String name = dir.toFile().getName(); + if (name.startsWith(DIRPREFIX)) { + String dateStr = name.substring(DIRPREFIX.length(),name.lastIndexOf(StringPool.UNDERSCORE)); + LocalDate dirDate = LocalDate.parse(dateStr, formatter); + if (!dirDate.isAfter(retentionDate)) { + targetDirs.add(dir); + } + } + } + } catch (IOException e) { + throw new RuntimeException(e); + } + }else { + LocalDate localDate = LocalDate.now().minusDays(gfsDataRetentionPeriod); + String batchTime = formatter.format(localDate); + String targetDir = this.gfsDataPath+ File.separator+DIRPREFIX+batchTime+StringPool.UNDERSCORE+HOURS; + targetDirs.add(Paths.get(targetDir)); + } + try { + for (Path targetDirPath : targetDirs){ + if(!Files.exists(targetDirPath)){ + break; + } + Files.walkFileTree(targetDirPath,new SimpleFileVisitor<>() { + @NotNull + @Override + public FileVisitResult visitFile(@NotNull Path file, @NotNull BasicFileAttributes attrs) throws IOException { + Files.delete(file); + return FileVisitResult.CONTINUE; + } + @NotNull + @Override + public FileVisitResult postVisitDirectory(@NotNull Path dir, @Nullable IOException exc) throws IOException { + Files.delete(dir); + return FileVisitResult.CONTINUE; + } + }); + } + this.cleanGFSDataLog(); + log.info("完成{}目录清除超出保留期限气象数据",gfsDataPath); + } catch (IOException e) { + log.error("{}目录清除超出保留期限气象数据出现异常",gfsDataPath,e); + } + } + + /** + * 清除gfs下载日志 + */ + private void cleanGFSDataLog(){ + LocalDateTime cleanDateTime = LocalDate.now().minusDays(gfsDataRetentionPeriod).atTime(23,59,59); + LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); + queryWrapper.le(WeatherDownGFSDataLog::getCreateTime,cleanDateTime); + this.downGFSDataLogMapper.delete(queryWrapper); + } + + /** + * 清除他t1h下载日志 + */ + private void cleanT1hDataLog(){ + LocalDateTime cleanDateTime = LocalDate.now().minusDays(gfsDataRetentionPeriod).atTime(23,59,59); + LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); + queryWrapper.le(WeatherDownT1hDataLog::getCreateTime,cleanDateTime); + this.downT1hDataLogMapper.delete(queryWrapper); + } +} diff --git a/jeecg-module-weather/src/main/java/org/jeecg/job/DownloadGFSJob.java b/jeecg-module-weather/src/main/java/org/jeecg/job/DownloadGFSJob.java new file mode 100644 index 0000000..b0071ea --- /dev/null +++ b/jeecg-module-weather/src/main/java/org/jeecg/job/DownloadGFSJob.java @@ -0,0 +1,274 @@ +package org.jeecg.job; + +import cn.hutool.core.collection.CollUtil; +import cn.hutool.core.util.StrUtil; +import com.baomidou.mybatisplus.core.toolkit.StringPool; +import com.jcraft.jsch.ChannelExec; +import com.jcraft.jsch.JSchException; +import com.jcraft.jsch.Session; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.codec.digest.DigestUtils; +import org.aspectj.util.FileUtil; +import org.jeecg.common.constant.enums.WeatherDataSourceEnum; +import org.jeecg.common.properties.ServerProperties; +import org.jeecg.common.util.JSchRemoteRunner; +import org.jeecg.modules.base.entity.WeatherData; +import org.jeecg.modules.base.entity.WeatherDownGFSDataLog; +import org.jeecg.modules.base.mapper.WeatherDataMapper; +import org.jeecg.modules.base.mapper.WeatherDownGFSDataLogMapper; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; +import java.io.*; +import java.time.LocalDate; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.concurrent.TimeUnit; + +/** + * 下载GFS预报数据任务 + */ +@Slf4j +@RequiredArgsConstructor +@Component +public class DownloadGFSJob { + + @Value("${gfs-download.downloadSwitch:false}") + private boolean downloadSwitch; + @Value("${gfs-download.pythonEnvPath}") + private String pythonEnvPath; + @Value("${gfs-download.gfsDownloadScriptPath}") + private String downloadScriptPath; + @Value("${gfs-download.gfsDataPath}") + private String gfsDataPath; + private final WeatherDownGFSDataLogMapper weatherDownGFSDataLogMapper; + private final ServerProperties serverProperties; + private final WeatherDataMapper weatherDataMapper; + + private final static String DATA_FILE_PREFIX = "gfs.t00z.pgrb2.0p25."; + private final static String HOURS = "00"; + private final static String DIRPREFIX = "GFS_"; + //下载失败重试次数 + private int retryCount = 0; + + @Scheduled(cron = "${gfs-download.cron}") + public void downloadT1HFile() { + //开关为true才执行下载任务 + if(downloadSwitch){ + retryCount = 0; + //ssh连接宿主机调用flexpart + JSchRemoteRunner jschRemoteRunner = new JSchRemoteRunner(); + try { + //登录ssh + jschRemoteRunner.login(this.serverProperties.getHost(),this.serverProperties.getPort(), + this.serverProperties.getUsername(),this.serverProperties.getPassword()); + //获取批次时间和下载命令 + DateTimeFormatter dateTimeFormatter = DateTimeFormatter.ofPattern("yyyyMMdd"); + String batchTime = dateTimeFormatter.format(LocalDate.now()); + this.execDownload(jschRemoteRunner,batchTime); + }catch (Exception e) { + log.error("GFS预报数据下载出现错误",e); + } + } + } + + /** + * 执行下载 + * @param jschRemoteRunner + */ + private void execDownload(JSchRemoteRunner jschRemoteRunner,String batchTime){ + //获取批次时间和下载命令 + String downloadCommand = this.getDownloadCommand(batchTime); + //执行下载命令 + log.info("开始下载"); + log.info("下载命令为{}",downloadCommand); + this.execCommand(batchTime,downloadCommand,jschRemoteRunner.getSession()); + //检查结果数据 + this.checkData(jschRemoteRunner,batchTime); + } + + /** + * 检验下载的数据是否完整 + * @param jschRemoteRunner + * @param batchTime + */ + private void checkData(JSchRemoteRunner jschRemoteRunner,String batchTime){ + log.info("开始检查数据"); + String gfsBatchDataPath = this.getGfsDataPath(batchTime); + List files = List.of(FileUtil.listFiles(new File(gfsBatchDataPath), file -> file.getName().startsWith(DATA_FILE_PREFIX))); + try { + if(CollUtil.isEmpty(files)){ + throw new RuntimeException(gfsBatchDataPath+"目录没有数据,下载可能出现错误"); + } + //根据业务规则生成预期的文件后缀集合 (f000 ~ f384) + Set expectedSuffixes = new HashSet<>(); + // 前 120 步进为 1 + for (int i = 0; i <= 120; i++) { + expectedSuffixes.add(String.format("f%03d", i)); + } + // 120 之后步进为 3 + for (int i = 123; i <= 384; i += 3) { + expectedSuffixes.add(String.format("f%03d", i)); + } + //获取实际后缀 + Set actualSuffixes = new HashSet<>(); + for (File file : files){ + actualSuffixes.add(file.getName().substring(DATA_FILE_PREFIX.length())); + } + // 比对预期与实际文件,找出缺失的文件 + boolean flag = false; + for (String expected : expectedSuffixes) { + if (!actualSuffixes.contains(expected)) { + String logContent = "缺失"+DATA_FILE_PREFIX+expected+"气象文件"; + this.saveLog(batchTime,logContent); + flag = true; + } + } + if (flag){ + throw new RuntimeException(gfsBatchDataPath+"目录批次数据存在缺失"); + } + //保存数据入库 + this.saveDataToDataBase(batchTime); + }catch (Exception e) { + this.saveLog(e.getMessage(),batchTime); + log.error(e.getMessage(),e); + try { + //间隔5分钟,重试200次 + if(retryCount <= 200){ + retryCount++; + TimeUnit.MINUTES.sleep(5); + log.info("{}批次数据,重试下载{}次,最大200次",gfsBatchDataPath,retryCount); + this.saveLog(batchTime,batchTime); + this.execDownload(jschRemoteRunner,batchTime); + } + } catch (InterruptedException ex) { + throw new RuntimeException(ex); + } + } + } + + /** + * 保存已下载的数据入库 + * @param batchTime + */ + private void saveDataToDataBase(String batchTime){ + //本批次数据开始日期 + LocalDateTime localDateTime = LocalDate.parse(batchTime,DateTimeFormatter.ofPattern("yyyyMMdd")).atTime(Integer.parseInt(HOURS),0,0); + String gfsBatchDataPath = this.getGfsDataPath(batchTime); + List fileList = List.of(FileUtil.listFiles(new File(gfsBatchDataPath), file -> file.getName().startsWith(DATA_FILE_PREFIX))); + if (CollUtil.isNotEmpty(fileList)) { + for (File file : fileList) { + try { + //获取文件数据时间 + int hour = Integer.parseInt(file.getName().substring(DATA_FILE_PREFIX.length()+1)); + LocalDateTime validTime = localDateTime.plusHours(hour); + //如果此文件存在则无需再次新增 + String gribFileMD5 = this.getGribFileMD5(file.getAbsolutePath()); + WeatherData weatherData = new WeatherData(); + weatherData.setFileName(file.getName()); + weatherData.setFileExt(file.getName().substring(file.getName().lastIndexOf(".")+1)); + weatherData.setDataSource(WeatherDataSourceEnum.GFS.getKey()); + weatherData.setFilePath(file.getAbsolutePath()); + weatherData.setDataStartTime(validTime); + weatherData.setMd5Value(gribFileMD5); + weatherData.setTimeBatch(batchTime); + this.weatherDataMapper.insert(weatherData); + }catch (Exception e){ + log.error("保存{}气象数据文件出现错误,原因为:",file.getAbsolutePath(),e); + + String logContent = "保存{}气象数据文件出现错误,原因为:"+e.getMessage(); + this.saveLog(logContent,batchTime); + } + } + } + } + + /** + * 获取GRIB文件的MD5唯一值 + */ + private String getGribFileMD5(String filePath) throws IOException { + try (FileInputStream fis = new FileInputStream(filePath)) { + // 底层自动采用流式读取,内存占用极低 + return DigestUtils.md5Hex(fis); + } + } + + /** + * 获取下载命令 + * @param batchTime + * @return + */ + private String getDownloadCommand(String batchTime){ + StringBuilder command = new StringBuilder(); + command.append(pythonEnvPath); + command.append(" -u "); + command.append(downloadScriptPath); + command.append(" --date "); + command.append(batchTime); + command.append(" --cycle "); + command.append(HOURS); + command.append(" --outdir "); + command.append(this.getGfsDataPath(batchTime)); + return command.toString(); + } + + /** + * 获取GFS数据存储路径 + * @param batchTime + * @return + */ + private String getGfsDataPath(String batchTime){ + return gfsDataPath + + File.separator + + DIRPREFIX + + batchTime + + StringPool.UNDERSCORE + + HOURS; + } + + /** + * 执行下载命令持续读取日志 + */ + protected void execCommand(String batchTime,String command, Session session){ + ChannelExec channel = null; + try{ + //打开一个执行通道 + channel = (ChannelExec) session.openChannel("exec"); + String fullCommand = command + " 2>&1"; + channel.setCommand(fullCommand); + // 获取脚本的标准输出流,包含错误输出流 + InputStream in = channel.getInputStream(); + BufferedReader reader = new BufferedReader(new InputStreamReader(in)); + // 连接通道 + channel.connect(); + + String line; + while ((line = reader.readLine()) != null) { + if(StrUtil.isNotBlank(line)){ + //保存日志 + this.saveLog(line,batchTime); + } + } + }catch(JSchException | IOException e){ + throw new RuntimeException(e); + }finally { + if (channel != null) { + channel.disconnect(); + } + } + } + + /** + * 保存过程日志 + */ + private void saveLog(String log,String batchTime){ + WeatherDownGFSDataLog dataLog = new WeatherDownGFSDataLog(); + dataLog.setLogContent(log); + dataLog.setBatchTime(batchTime); + weatherDownGFSDataLogMapper.insert(dataLog); + } +} 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 fb47bc0..26428ff 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 @@ -75,6 +75,8 @@ public class WeatherDataServiceImpl extends ServiceImpl