|
|
@@ -12,6 +12,7 @@ import java.util.HashMap;
|
|
|
import java.util.LinkedHashMap;
|
|
|
import java.util.List;
|
|
|
import java.util.Map;
|
|
|
+import java.util.TreeMap;
|
|
|
import java.util.function.Consumer;
|
|
|
import java.util.function.Function;
|
|
|
import java.util.function.Supplier;
|
|
|
@@ -34,7 +35,9 @@ import com.ufo.project.forecast.mediumterm.power.mapper.DataMediumTermPowerForec
|
|
|
import com.ufo.project.forecast.shortterm.power.domain.DataShortTermPowerForecastRaw;
|
|
|
import com.ufo.project.forecast.shortterm.power.mapper.DataShortTermPowerForecastRawMapper;
|
|
|
import com.ufo.project.forecast.supershort.domain.DataSuperShortTermForecastRaw;
|
|
|
+import com.ufo.project.forecast.supershort.domain.DataUltraShortTermPowerForecasting;
|
|
|
import com.ufo.project.forecast.supershort.mapper.DataSuperShortTermForecastRawMapper;
|
|
|
+import com.ufo.project.forecast.supershort.mapper.DataUltraShortTermPowerForecastingMapper;
|
|
|
|
|
|
/**
|
|
|
* 功率预测数据文件处理器(短期/中期/超短期)
|
|
|
@@ -100,6 +103,12 @@ public class DataFileProcesserPowerImpl extends AbstraceDataFileProcesser
|
|
|
private static final Long ANALYZE_HIS_ID = 0L;
|
|
|
private static final Long ARCHIVE_RECORD_ID = 0L;
|
|
|
|
|
|
+ /** data_ultra_short_term_power_forecasting 的 powerXX 列: 分钟间隔步长与最大时效(t+15..t+240) */
|
|
|
+ private static final long FORECAST_STEP_MINUTES = 15L;
|
|
|
+ private static final long FORECAST_MAX_MINUTES = 240L;
|
|
|
+ /** data_ultra_short_term_power_forecasting 的 entity_type(场站/交易主体/省域), 场站取1 */
|
|
|
+ private static final Long ULTRA_ENTITY_TYPE = 1L;
|
|
|
+
|
|
|
/** 短期/中期/超短期各自的元数据组(id同时用作data_source/data_sourcegrab_id/data_model_id/source_id/subject_id/machine_id) */
|
|
|
static final Term SHORT_TERM = new Term(Kind.SHORT, 5001L, "EC", "ShortTermModel", "短期", 3L, 72L);
|
|
|
static final Term MEDIUM_TERM = new Term(Kind.MEDIUM, 5002L, "WRF", "MediumTermModel", "中期", 10L, 240L);
|
|
|
@@ -114,6 +123,9 @@ public class DataFileProcesserPowerImpl extends AbstraceDataFileProcesser
|
|
|
@Autowired
|
|
|
private DataSuperShortTermForecastRawMapper dataSuperShortTermForecastRawMapper;
|
|
|
|
|
|
+ @Autowired
|
|
|
+ private DataUltraShortTermPowerForecastingMapper dataUltraShortTermPowerForecastingMapper;
|
|
|
+
|
|
|
@Override
|
|
|
protected String parseAndSave(DataFileSyncFiles file)
|
|
|
{
|
|
|
@@ -153,7 +165,7 @@ public class DataFileProcesserPowerImpl extends AbstraceDataFileProcesser
|
|
|
return "csv无有效数据行";
|
|
|
}
|
|
|
|
|
|
- int[] result;
|
|
|
+ int[] result = new int[]{0, 0};
|
|
|
switch (term.kind)
|
|
|
{
|
|
|
case SHORT:
|
|
|
@@ -162,8 +174,11 @@ public class DataFileProcesserPowerImpl extends AbstraceDataFileProcesser
|
|
|
case MEDIUM:
|
|
|
result = saveMediumTerm(term, station, startTime, powerByForecastTime);
|
|
|
break;
|
|
|
- default:
|
|
|
+ case ULTRA:
|
|
|
result = saveUltraShortTerm(term, station, startTime, powerByForecastTime);
|
|
|
+ saveUltraShortTermNew(term, station, startTime, powerByForecastTime);
|
|
|
+ break;
|
|
|
+ default:
|
|
|
break;
|
|
|
}
|
|
|
log.info("[DataFileProcesserPower] 处理完成, file={}, 场站={}, {}, start_time={}, 更新{}条, 新增{}条",
|
|
|
@@ -336,6 +351,65 @@ public class DataFileProcesserPowerImpl extends AbstraceDataFileProcesser
|
|
|
dataSuperShortTermForecastRawMapper::batchInsertDataSuperShortTermForecastRaw);
|
|
|
}
|
|
|
|
|
|
+ /**
|
|
|
+ * 合并保存一条 data_ultra_short_term_power_forecasting: 按(entity_id, start_time)查重,
|
|
|
+ * 已存在则按主键覆盖powerXX预测列(不动实发power), 否则新增(create_time取库端当前epoch秒)
|
|
|
+ */
|
|
|
+ private void saveUltraShortTermNew(Term term, ConfigStation station, long startTime, Map<Long, String> powerByForecastTime)
|
|
|
+ {
|
|
|
+ DataUltraShortTermPowerForecasting row = buildUltraForecasting(term, station, startTime, powerByForecastTime);
|
|
|
+ DataUltraShortTermPowerForecasting query = new DataUltraShortTermPowerForecasting();
|
|
|
+ query.setEntityId(row.getEntityId());
|
|
|
+ query.setStartTime(startTime);
|
|
|
+ List<DataUltraShortTermPowerForecasting> dbList =
|
|
|
+ dataUltraShortTermPowerForecastingMapper.selectDataUltraShortTermPowerForecastingList(query);
|
|
|
+ if (dbList != null && !dbList.isEmpty() && dbList.get(0).getId() != null)
|
|
|
+ {
|
|
|
+ row.setId(dbList.get(0).getId());
|
|
|
+ dataUltraShortTermPowerForecastingMapper.updateDataUltraShortTermPowerForecastingFromFile(row);
|
|
|
+ }
|
|
|
+ else
|
|
|
+ {
|
|
|
+ dataUltraShortTermPowerForecastingMapper.insertDataUltraShortTermPowerForecastingFromFile(row);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 组装超短期透视行: powerByForecastTime按key升序, 以与startTime的间隔分钟数取对应powerXX列
|
|
|
+ * (仅t+15..t+240的15分钟整倍数, 其余跳过并告警); 元数据与raw链路同源(term=5003/EC-WRF),
|
|
|
+ * entity_id为场站id字符串、entity_type场站取1、entity_name场站名;
|
|
|
+ * forecast_time取最后一个有效预测时间(列可空), 实发power不写
|
|
|
+ */
|
|
|
+ DataUltraShortTermPowerForecasting buildUltraForecasting(Term term, ConfigStation station, long startTime, Map<Long, String> powerByForecastTime)
|
|
|
+ {
|
|
|
+ DataUltraShortTermPowerForecasting row = new DataUltraShortTermPowerForecasting();
|
|
|
+ row.setDataSourcegrabId(term.id);
|
|
|
+ row.setDataSourcegrabPatternName(term.patternName);
|
|
|
+ row.setEntityId(String.valueOf(station.getId()));
|
|
|
+ row.setEntityType(ULTRA_ENTITY_TYPE);
|
|
|
+ row.setEntityName(station.getStationName());
|
|
|
+ row.setStartTime(startTime);
|
|
|
+ Long lastForecastTime = null;
|
|
|
+ for (Map.Entry<Long, String> entry : new TreeMap<>(powerByForecastTime).entrySet())
|
|
|
+ {
|
|
|
+ long offsetSeconds = entry.getKey() - startTime;
|
|
|
+ long minutes = offsetSeconds / 60L;
|
|
|
+ if (offsetSeconds > 0 && offsetSeconds % (FORECAST_STEP_MINUTES * 60L) == 0
|
|
|
+ && minutes >= FORECAST_STEP_MINUTES && minutes <= FORECAST_MAX_MINUTES)
|
|
|
+ {
|
|
|
+ setEntityField(row, "power" + minutes, entry.getValue());
|
|
|
+ lastForecastTime = entry.getKey();
|
|
|
+ }
|
|
|
+ else
|
|
|
+ {
|
|
|
+ log.warn("[DataFileProcesserPower] 超短期预测时间与起报间隔不符合powerXX列, 跳过, start_time={}, forecast_time={}",
|
|
|
+ startTime, entry.getKey());
|
|
|
+ }
|
|
|
+ }
|
|
|
+ row.setForecastTime(lastForecastTime);
|
|
|
+ return row;
|
|
|
+ }
|
|
|
+
|
|
|
/**
|
|
|
* 三条链路共用: 按(station_id, start_time)查已有记录, 以forecast_time匹配,
|
|
|
* 已有置Id进更新组、其余进插入组, 各按BATCH_SIZE一批写库;
|