feat(application): integrate measured-load ingestion training and planner source
This commit is contained in:
@@ -0,0 +1,908 @@
|
||||
from netplan_v4_publisher import publish_forecasts as _v4_publish_forecasts
|
||||
from model_isolation import collect_predictions
|
||||
import os
|
||||
import sqlite3
|
||||
import json
|
||||
import time
|
||||
import datetime
|
||||
import traceback
|
||||
import warnings
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
|
||||
import numpy as np
|
||||
import pandas as pd
|
||||
import pytz
|
||||
from influxdb_client import InfluxDBClient, Point, WritePrecision
|
||||
from influxdb_client.client.write_api import SYNCHRONOUS
|
||||
from soc_diagnostics import battery_soc_points
|
||||
from telemetry_quality import require_recent_telemetry, sanitize_measured_frame
|
||||
try:
|
||||
from influxdb_client.client.warnings import MissingPivotFunction
|
||||
warnings.simplefilter("ignore", MissingPivotFunction)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
import methods.var_1 as v1
|
||||
import methods.var_2 as v2
|
||||
import methods.var_3 as v3
|
||||
import methods.var_10 as v10
|
||||
import methods.var_11 as v11
|
||||
import methods.var_13 as v13
|
||||
import methods.var_21 as v21
|
||||
import methods.var_22 as v22
|
||||
import methods.var_23 as v23
|
||||
|
||||
INFLUX_URL = os.getenv("INFLUX_URL", "http://influxdb:8086")
|
||||
INFLUX_TOKEN = os.environ["INFLUX_TOKEN"]
|
||||
INFLUX_ORG = os.getenv("INFLUX_ORG", "belevo")
|
||||
INFLUX_BUCKET = os.getenv("INFLUX_BUCKET", "energy_data")
|
||||
SQLITE_DB_PATH = os.getenv("SQLITE_DB_PATH", "/app/data/users.db")
|
||||
HISTORY_START = os.getenv("FORECAST_HISTORY_START", "1970-01-01T00:00:00Z")
|
||||
QUALITY_LOOKBACK_DAYS = int(os.getenv("FORECAST_QUALITY_LOOKBACK_DAYS", "14"))
|
||||
LOCAL_TZ = pytz.timezone(os.getenv("TZ", "Europe/Zurich"))
|
||||
INFLUX_TIMEOUT_MS = int(os.getenv("FORECAST_INFLUX_TIMEOUT_MS", "120000"))
|
||||
FORECAST_TIMEOUT_SECONDS = int(os.getenv("FORECAST_RUN_TIMEOUT_SECONDS", "900"))
|
||||
FORECAST_HORIZON_HOURS = max(24, min(72, int(os.getenv("FORECAST_HORIZON_HOURS", "48"))))
|
||||
TRAINING_TIMEOUT_SECONDS = int(os.getenv("FORECAST_TRAINING_TIMEOUT_SECONDS", "3600"))
|
||||
FORECAST_STALE_SECONDS = int(os.getenv("FORECAST_STALE_SECONDS", "5400"))
|
||||
WATCHDOG_INTERVAL_SECONDS = int(os.getenv("FORECAST_WATCHDOG_INTERVAL_SECONDS", "300"))
|
||||
TRAINING_HOUR = int(os.getenv("FORECAST_TRAINING_HOUR", "2"))
|
||||
LAST_SUCCESS_PATH = os.getenv("FORECAST_LAST_SUCCESS_PATH", "/tmp/forecast_engine_last_success.json")
|
||||
LAST_TRAINING_PATH = os.getenv("FORECAST_LAST_TRAINING_PATH", "/app/data/forecast_training_status.json")
|
||||
last_trained_day = None
|
||||
_FORECAST_PROCESS_LOCK = threading.Lock()
|
||||
_TRAINING_THREAD = None
|
||||
|
||||
MODEL_MODULES = {1: v1, 2: v2, 3: v3, 10: v10, 11: v11, 13: v13, 21: v21, 22: v22, 23: v23}
|
||||
QUALITY_TARGETS = {1: "PV", 10: "PV", 21: "PV", 2: "Hausverbrauch", 11: "Hausverbrauch", 22: "Hausverbrauch"}
|
||||
|
||||
|
||||
def active(config, n):
|
||||
return bool(int(config.get(f"prog_var_{n}", config.get(f"var_{n}", 0)) or 0))
|
||||
|
||||
|
||||
def get_configs():
|
||||
conn = sqlite3.connect(SQLITE_DB_PATH)
|
||||
conn.row_factory = sqlite3.Row
|
||||
columns = {row["name"] for row in conn.execute("PRAGMA table_info(anlagen_meta)")}
|
||||
if "batt_grid_charging_enabled" not in columns:
|
||||
conn.execute(
|
||||
"ALTER TABLE anlagen_meta ADD COLUMN batt_grid_charging_enabled INTEGER NOT NULL DEFAULT 0"
|
||||
)
|
||||
conn.commit()
|
||||
rows = conn.execute("SELECT * FROM anlagen_meta").fetchall()
|
||||
conn.close()
|
||||
configs = []
|
||||
for r in rows:
|
||||
d = dict(r)
|
||||
try:
|
||||
d["daecher"] = json.loads(d.get("daecher") or "[]")
|
||||
except Exception:
|
||||
d["daecher"] = []
|
||||
for key, default in [
|
||||
("ac_leistung", 10.0),
|
||||
("batt_capacity_kwh", 0.0),
|
||||
("batt_power_kw", 0.0),
|
||||
("tarif_bezug_fest", 0.30),
|
||||
("tarif_einspeisung_fest", 0.10),
|
||||
("tarif_peak_fest", 5.0),
|
||||
]:
|
||||
raw_value = d.get(key)
|
||||
d[key] = float(default if raw_value is None or raw_value == "" else raw_value)
|
||||
configs.append(d)
|
||||
return configs
|
||||
|
||||
|
||||
def _time_literal(dt):
|
||||
return dt.strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
|
||||
|
||||
def _query_df(query):
|
||||
client = InfluxDBClient(
|
||||
url=INFLUX_URL,
|
||||
token=INFLUX_TOKEN,
|
||||
org=INFLUX_ORG,
|
||||
timeout=INFLUX_TIMEOUT_MS,
|
||||
)
|
||||
try:
|
||||
df = client.query_api().query_data_frame(org=INFLUX_ORG, query=query)
|
||||
finally:
|
||||
client.close()
|
||||
if isinstance(df, list):
|
||||
df = pd.concat(df, ignore_index=True) if df else pd.DataFrame()
|
||||
if df is None or df.empty or "_time" not in df.columns:
|
||||
return pd.DataFrame()
|
||||
df["_time"] = pd.to_datetime(df["_time"], utc=True).dt.tz_localize(None)
|
||||
return df
|
||||
|
||||
|
||||
def _pivot_frame(df, fields):
|
||||
if df.empty:
|
||||
return pd.DataFrame()
|
||||
out = df.set_index("_time")
|
||||
keep = [c for c in fields if c in out.columns]
|
||||
out = out[keep] if keep else pd.DataFrame(index=out.index)
|
||||
out = out.apply(pd.to_numeric, errors="coerce")
|
||||
out = out[~out.index.duplicated(keep="last")]
|
||||
return out.sort_index().resample("5min").mean(numeric_only=True)
|
||||
|
||||
|
||||
def _tariff_frame(df, config=None):
|
||||
if df is None or df.empty:
|
||||
return pd.DataFrame()
|
||||
df = df.copy()
|
||||
if "_time" in df.columns:
|
||||
df["_time"] = pd.to_datetime(df["_time"], errors="coerce")
|
||||
df = df.dropna(subset=["_time"]).set_index("_time")
|
||||
if df.empty:
|
||||
return pd.DataFrame()
|
||||
|
||||
price = df["price_chf_kwh"] if "price_chf_kwh" in df.columns else df.get("_value")
|
||||
if price is None:
|
||||
return pd.DataFrame()
|
||||
price = pd.to_numeric(price, errors="coerce")
|
||||
model = df.get("tariff_model", pd.Series("", index=df.index)).astype(str).str.lower()
|
||||
typ = df.get("type", pd.Series("", index=df.index)).astype(str).str.lower()
|
||||
tariff_name = df.get("tariff_name", pd.Series("", index=df.index)).astype(str).str.lower()
|
||||
provider = df.get("provider", pd.Series("", index=df.index)).astype(str).str.lower()
|
||||
key = (model + " " + tariff_name + " " + provider).str.lower()
|
||||
|
||||
rows = pd.DataFrame({"price": price, "key": key, "type": typ}, index=df.index)
|
||||
rows = rows[pd.notna(rows["price"])].sort_index()
|
||||
if rows.empty:
|
||||
return pd.DataFrame()
|
||||
|
||||
cfg = config or {}
|
||||
import_choice = str(cfg.get("tarif_bezug", "") or "").lower()
|
||||
export_choice = str(cfg.get("tarif_einspeisung", "") or "").lower()
|
||||
|
||||
out = pd.DataFrame(index=rows.index.unique().sort_values())
|
||||
|
||||
import_base = (
|
||||
rows["type"].str.contains("consumption|import|bezug", regex=True, na=False)
|
||||
| rows["key"].str.contains("dynamic|dynamisch|home|business", regex=True, na=False)
|
||||
)
|
||||
if "business" in import_choice:
|
||||
import_mask = import_base & rows["key"].str.contains("business|gewerbe|commercial", regex=True, na=False)
|
||||
elif "home" in import_choice or "privat" in import_choice:
|
||||
import_mask = import_base & rows["key"].str.contains("home|privat|private", regex=True, na=False)
|
||||
elif "dynam" in import_choice:
|
||||
import_mask = import_base
|
||||
else:
|
||||
import_mask = pd.Series(False, index=rows.index)
|
||||
if not import_mask.any() and "dynam" in import_choice:
|
||||
import_mask = import_base
|
||||
|
||||
export_base = (
|
||||
rows["type"].str.contains("feed|einspeis|export", regex=True, na=False)
|
||||
| rows["key"].str.contains("referenzmarktpreis|marktpreis|reference|feed", regex=True, na=False)
|
||||
)
|
||||
if "referenz" in export_choice or "marktpreis" in export_choice or "market" in export_choice:
|
||||
export_mask = export_base & rows["key"].str.contains("referenzmarktpreis|marktpreis|reference|belevo", regex=True, na=False)
|
||||
else:
|
||||
export_mask = rows["key"].str.contains("standard_feedin|ckw statisch", regex=True, na=False) & export_base
|
||||
if not export_mask.any() and ("referenz" in export_choice or "marktpreis" in export_choice or "market" in export_choice):
|
||||
export_mask = export_base
|
||||
|
||||
if import_mask.any():
|
||||
out["import_price"] = rows.loc[import_mask, "price"].groupby(level=0).last()
|
||||
if export_mask.any():
|
||||
out["export_price"] = rows.loc[export_mask, "price"].groupby(level=0).last()
|
||||
|
||||
if out.empty:
|
||||
return out
|
||||
return out.sort_index().resample("5min").mean().ffill().bfill()
|
||||
|
||||
def fetch_influx_frames(config, training):
|
||||
aid = config["anlagen_id"]
|
||||
start = HISTORY_START if training else "-14d"
|
||||
future_stop = _time_literal(datetime.datetime.utcnow() + datetime.timedelta(hours=FORECAST_HORIZON_HOURS))
|
||||
|
||||
q_tel = f'''
|
||||
from(bucket: "{INFLUX_BUCKET}")
|
||||
|> range(start: {start})
|
||||
|> filter(fn: (r) => r["_measurement"] == "api_telemetry")
|
||||
|> filter(fn: (r) => r["anlagen_id"] == "{aid}")
|
||||
|> filter(fn: (r) => r["_field"] == "PV" or r["_field"] == "Hausverbrauch" or r["_field"] == "Netzleistung" or r["_field"] == "SOC")
|
||||
|> filter(fn: (r) => not exists r["data_type"] or (r["data_type"] != "forecast" and r["data_type"] != "forecast_snapshot"))
|
||||
|> aggregateWindow(every: 5m, fn: mean, createEmpty: false, timeSrc: "_start")
|
||||
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
||||
'''
|
||||
df_tel = _pivot_frame(_query_df(q_tel), ["PV", "Hausverbrauch", "Netzleistung", "SOC"])
|
||||
|
||||
q_wea = f'''
|
||||
from(bucket: "{INFLUX_BUCKET}")
|
||||
|> range(start: {start}, stop: {future_stop})
|
||||
|> filter(fn: (r) => r["_measurement"] == "weather_forecast")
|
||||
|> aggregateWindow(every: 5m, fn: mean, createEmpty: false)
|
||||
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
||||
'''
|
||||
df_wea = _pivot_frame(_query_df(q_wea), ["temp_c", "temperature", "cloud", "cloud_cover", "precip_mm", "wind_kph", "chance_of_snow"])
|
||||
if "temperature" in df_wea.columns and "temp_c" not in df_wea.columns:
|
||||
df_wea["temp_c"] = df_wea["temperature"]
|
||||
if "cloud_cover" in df_wea.columns and "cloud" not in df_wea.columns:
|
||||
df_wea["cloud"] = df_wea["cloud_cover"]
|
||||
if "cloud" in df_wea.columns and "cloud_cover" not in df_wea.columns:
|
||||
df_wea["cloud_cover"] = df_wea["cloud"]
|
||||
|
||||
q_tar = f'''
|
||||
from(bucket: "{INFLUX_BUCKET}")
|
||||
|> range(start: -7d, stop: {future_stop})
|
||||
|> filter(fn: (r) => r["_measurement"] == "tariffs")
|
||||
|> filter(fn: (r) => r["_field"] == "price_chf_kwh")
|
||||
'''
|
||||
df_tar = _tariff_frame(_query_df(q_tar), config)
|
||||
return df_tel, df_wea, df_tar
|
||||
|
||||
|
||||
def _consistent_tail(df):
|
||||
if df.empty or "PV" not in df.columns or "Hausverbrauch" not in df.columns:
|
||||
return df
|
||||
probe = df[["PV", "Hausverbrauch"]].copy()
|
||||
filled = probe.interpolate(limit=3, limit_direction="both")
|
||||
valid = filled.notna().all(axis=1)
|
||||
if not valid.any():
|
||||
return df.iloc[0:0]
|
||||
run = 0
|
||||
last_break = -1
|
||||
for i, ok in enumerate(valid.to_numpy()):
|
||||
if ok:
|
||||
run = 0
|
||||
else:
|
||||
run += 1
|
||||
if run >= 4:
|
||||
last_break = i
|
||||
if last_break >= 0:
|
||||
after = np.where(valid.iloc[last_break + 1:].to_numpy())[0]
|
||||
if len(after):
|
||||
return df.iloc[last_break + 1 + after[0]:]
|
||||
return df.iloc[0:0]
|
||||
return df.loc[valid[valid].index[0]:]
|
||||
|
||||
|
||||
|
||||
def _longest_consistent_segment(df, column):
|
||||
if df.empty or column not in df.columns:
|
||||
return pd.DataFrame()
|
||||
series = pd.to_numeric(df[column], errors="coerce")
|
||||
valid = series.notna().to_numpy()
|
||||
if not valid.any():
|
||||
return pd.DataFrame()
|
||||
best_start = best_end = None
|
||||
start = 0
|
||||
gap = 0
|
||||
for i, ok in enumerate(valid):
|
||||
if ok:
|
||||
gap = 0
|
||||
else:
|
||||
gap += 1
|
||||
if gap >= 4:
|
||||
end = i - gap
|
||||
if end >= start and (best_start is None or end - start > best_end - best_start):
|
||||
best_start, best_end = start, end
|
||||
start = i + 1
|
||||
gap = 0
|
||||
end = len(valid) - 1
|
||||
if end >= start and (best_start is None or end - start > best_end - best_start):
|
||||
best_start, best_end = start, end
|
||||
if best_start is None:
|
||||
return pd.DataFrame()
|
||||
segment = df.iloc[best_start:best_end + 1].copy()
|
||||
first = pd.to_numeric(segment[column], errors="coerce").first_valid_index()
|
||||
last = pd.to_numeric(segment[column], errors="coerce").last_valid_index()
|
||||
if first is None or last is None:
|
||||
return pd.DataFrame()
|
||||
return segment.loc[first:last]
|
||||
|
||||
|
||||
def _add_time_features(df):
|
||||
hour = df.index.hour + df.index.minute / 60.0
|
||||
doy = df.index.dayofyear
|
||||
df["hour_float"] = hour
|
||||
df["hour_sin"] = np.sin(2 * np.pi * hour / 24.0)
|
||||
df["hour_cos"] = np.cos(2 * np.pi * hour / 24.0)
|
||||
df["sin_year"] = np.sin(2 * np.pi * doy / 365.25)
|
||||
df["cos_year"] = np.cos(2 * np.pi * doy / 365.25)
|
||||
df["weekday"] = df.index.weekday
|
||||
df["is_weekday"] = (df.index.weekday < 5).astype(int)
|
||||
return df
|
||||
|
||||
|
||||
def _fill_defaults(df, history):
|
||||
defaults = {
|
||||
"PV": np.nan,
|
||||
"Hausverbrauch": np.nan,
|
||||
"SOC": np.nan,
|
||||
"Netzleistung": np.nan,
|
||||
"temp_c": 15.0,
|
||||
"cloud": 20.0,
|
||||
"cloud_cover": 20.0,
|
||||
"precip_mm": 0.0,
|
||||
"wind_kph": 0.0,
|
||||
"chance_of_snow": 0.0,
|
||||
"import_price": np.nan,
|
||||
"export_price": np.nan,
|
||||
}
|
||||
for col, default in defaults.items():
|
||||
if col not in df.columns:
|
||||
df[col] = default
|
||||
df[col] = pd.to_numeric(df[col], errors="coerce")
|
||||
# Measurements are never interpolated/filled here. Outages are not zero load,
|
||||
# zero PV, a fresh SOC, or a measured grid peak. Weather defaults are separate.
|
||||
df = sanitize_measured_frame(df)
|
||||
for col, default in defaults.items():
|
||||
if col not in ["PV", "Hausverbrauch", "SOC", "Netzleistung"]:
|
||||
df[col] = df[col].ffill().bfill()
|
||||
if np.isfinite(default):
|
||||
df[col] = df[col].fillna(default)
|
||||
if "cloud_cover" in df.columns:
|
||||
df["cloud"] = df["cloud"].fillna(df["cloud_cover"])
|
||||
df["import_price"] = df["import_price"].fillna(0.30)
|
||||
df["export_price"] = df["export_price"].fillna(0.10)
|
||||
return _add_time_features(df)
|
||||
|
||||
|
||||
def build_data_object(config, training=False):
|
||||
now = datetime.datetime.utcnow().replace(second=0, microsecond=0)
|
||||
now = now - datetime.timedelta(minutes=now.minute % 5)
|
||||
df_tel, df_wea, df_tar = fetch_influx_frames(config, training)
|
||||
|
||||
frames = [f for f in [df_tel, df_wea, df_tar] if not f.empty]
|
||||
combined = frames[0] if frames else pd.DataFrame()
|
||||
for f in frames[1:]:
|
||||
combined = combined.join(f, how="outer")
|
||||
combined = combined.sort_index()
|
||||
|
||||
full_hist_raw = combined.loc[:now - datetime.timedelta(minutes=5)] if not combined.empty else pd.DataFrame()
|
||||
# A trailing telemetry gap must not discard all earlier valid observations.
|
||||
hist_raw = full_hist_raw.copy()
|
||||
# Freshness is determined from measurements, not the outer-joined weather grid.
|
||||
recent_raw = sanitize_measured_frame(df_tel.loc[:now - datetime.timedelta(minutes=5)].copy()) if not df_tel.empty else pd.DataFrame()
|
||||
start_hist = hist_raw.index.min().floor("5min") if training and not hist_raw.empty else now - datetime.timedelta(days=14)
|
||||
|
||||
idx_hist = pd.date_range(start=start_hist, end=now - datetime.timedelta(minutes=5), freq="5min")
|
||||
df_hist = pd.DataFrame(index=idx_hist).join(hist_raw, how="left")
|
||||
df_hist = _fill_defaults(df_hist, history=True)
|
||||
|
||||
df_load_training = pd.DataFrame()
|
||||
if training and "Hausverbrauch" in full_hist_raw.columns:
|
||||
load_segment = _longest_consistent_segment(full_hist_raw, "Hausverbrauch")
|
||||
if not load_segment.empty:
|
||||
idx_load = pd.date_range(
|
||||
start=load_segment.index.min().floor("5min"),
|
||||
end=load_segment.index.max().floor("5min"),
|
||||
freq="5min",
|
||||
)
|
||||
df_load_training = pd.DataFrame(index=idx_load).join(load_segment, how="left")
|
||||
df_load_training = _fill_defaults(df_load_training, history=True)
|
||||
|
||||
df_pv_training = pd.DataFrame()
|
||||
if training and "PV" in full_hist_raw.columns:
|
||||
pv_first = full_hist_raw["PV"].first_valid_index()
|
||||
if pv_first is not None:
|
||||
idx_pv = pd.date_range(
|
||||
start=pv_first.floor("5min"),
|
||||
end=now - datetime.timedelta(minutes=5),
|
||||
freq="5min",
|
||||
)
|
||||
df_pv_training = pd.DataFrame(index=idx_pv).join(full_hist_raw, how="left")
|
||||
df_pv_training = _fill_defaults(df_pv_training, history=True)
|
||||
|
||||
idx_fut = pd.date_range(start=now, periods=FORECAST_HORIZON_HOURS * 12, freq="5min")
|
||||
fut_raw = combined.reindex(combined.index.union(idx_fut)).sort_index() if not combined.empty else pd.DataFrame(index=idx_fut)
|
||||
df_fut = pd.DataFrame(index=idx_fut).join(fut_raw, how="left")
|
||||
df_fut = _fill_defaults(df_fut, history=False)
|
||||
|
||||
month_hist = df_hist[(df_hist.index.year == now.year) & (df_hist.index.month == now.month)]
|
||||
current_peak_kw = 0.0
|
||||
if not month_hist.empty and "Netzleistung" in month_hist.columns:
|
||||
measured_grid = pd.to_numeric(month_hist["Netzleistung"], errors="coerce").dropna()
|
||||
measured_peak = measured_grid.resample(
|
||||
"15min", origin="start_day", label="left", closed="left"
|
||||
).mean()
|
||||
if not measured_peak.empty:
|
||||
current_peak_kw = max(0.0, float(measured_peak.max()) / 1000.0)
|
||||
if current_peak_kw <= 0.0 and not month_hist.empty and {"Hausverbrauch", "PV"}.issubset(month_hist.columns):
|
||||
fallback_residual = (month_hist["Hausverbrauch"] - month_hist["PV"]).resample(
|
||||
"15min", origin="start_day", label="left", closed="left"
|
||||
).mean()
|
||||
if not fallback_residual.empty:
|
||||
current_peak_kw = max(0.0, float(fallback_residual.max()) / 1000.0)
|
||||
|
||||
min_soc = float(
|
||||
config.get("batt_min_soc", config.get("batt_min_soc_percent", 0.0)) or 0.0
|
||||
)
|
||||
current_soc = max(0.0, min(100.0, min_soc))
|
||||
current_soc_source = "safe_minimum"
|
||||
current_soc_age_minutes = None
|
||||
if "SOC" in full_hist_raw.columns:
|
||||
soc_values = pd.to_numeric(full_hist_raw["SOC"], errors="coerce").dropna()
|
||||
if not soc_values.empty:
|
||||
latest_soc_time = pd.Timestamp(soc_values.index[-1])
|
||||
current_soc_age_minutes = max(
|
||||
0.0,
|
||||
(pd.Timestamp(now) - latest_soc_time).total_seconds() / 60.0,
|
||||
)
|
||||
try:
|
||||
max_soc_age_minutes = max(
|
||||
5.0,
|
||||
float(config.get("batt_soc_max_age_minutes", 30.0) or 30.0),
|
||||
)
|
||||
except (TypeError, ValueError):
|
||||
max_soc_age_minutes = 30.0
|
||||
if current_soc_age_minutes <= max_soc_age_minutes:
|
||||
current_soc = float(soc_values.iloc[-1])
|
||||
current_soc_source = "telemetry"
|
||||
current_soc = max(0.0, min(100.0, current_soc))
|
||||
|
||||
return {
|
||||
"config": config,
|
||||
"metadata": config,
|
||||
"now": now,
|
||||
"df_hist": df_hist,
|
||||
"df_recent_raw": recent_raw,
|
||||
"df_load_training": df_load_training,
|
||||
"df_pv_training": df_pv_training,
|
||||
"df_fut": df_fut,
|
||||
"current_soc_perc": current_soc,
|
||||
"current_soc_source": current_soc_source,
|
||||
"current_soc_age_minutes": current_soc_age_minutes,
|
||||
"current_month_peak_kw": current_peak_kw,
|
||||
}
|
||||
|
||||
|
||||
def run_training():
|
||||
for config in get_configs():
|
||||
if not any(active(config, n) for n in MODEL_MODULES):
|
||||
continue
|
||||
print(f"Training Anlage {config['anlagen_id']}...")
|
||||
data_obj = build_data_object(config, training=True)
|
||||
for n, module in MODEL_MODULES.items():
|
||||
if active(config, n) and hasattr(module, "train"):
|
||||
try:
|
||||
target = QUALITY_TARGETS.get(n)
|
||||
if target:
|
||||
require_recent_telemetry(data_obj, [target])
|
||||
print(f" var_{n}: {module.train(data_obj)}")
|
||||
except Exception:
|
||||
print(f" var_{n}: Training fehlgeschlagen")
|
||||
traceback.print_exc()
|
||||
|
||||
|
||||
def _forecast_point(aid, field, t, value):
|
||||
return (
|
||||
Point("api_telemetry")
|
||||
.tag("anlagen_id", aid)
|
||||
.tag("data_type", "forecast")
|
||||
.field(field, float(value))
|
||||
.time(t.to_pydatetime(), WritePrecision.S)
|
||||
)
|
||||
|
||||
|
||||
def _snapshot_point(aid, field, run_hour, t, value):
|
||||
return (
|
||||
Point("api_telemetry")
|
||||
.tag("anlagen_id", aid)
|
||||
.tag("data_type", "forecast_snapshot")
|
||||
.tag("run_hour", run_hour)
|
||||
.field(f"{field}_run_{run_hour}", float(value))
|
||||
.time(t.to_pydatetime(), WritePrecision.S)
|
||||
)
|
||||
|
||||
|
||||
def _quality_query(aid, forecast_field, target_field):
|
||||
stop = _time_literal(datetime.datetime.utcnow() - datetime.timedelta(minutes=10))
|
||||
return f'''
|
||||
from(bucket: "{INFLUX_BUCKET}")
|
||||
|> range(start: -{QUALITY_LOOKBACK_DAYS}d, stop: {stop})
|
||||
|> filter(fn: (r) => r["_measurement"] == "api_telemetry")
|
||||
|> filter(fn: (r) => r["anlagen_id"] == "{aid}")
|
||||
|> filter(fn: (r) => r["_field"] == "{target_field}" or r["_field"] == "{forecast_field}")
|
||||
|> aggregateWindow(every: 5m, fn: mean, createEmpty: false)
|
||||
|> group()
|
||||
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
||||
'''
|
||||
|
||||
|
||||
def _forecast_quality(aid, variant, target_field):
|
||||
forecast_field = f"prog_var_{variant}"
|
||||
df = _query_df(_quality_query(aid, forecast_field, target_field))
|
||||
if df.empty or forecast_field not in df.columns or target_field not in df.columns:
|
||||
return None
|
||||
pair = df[[target_field, forecast_field]].apply(pd.to_numeric, errors="coerce").dropna()
|
||||
if len(pair) < 3:
|
||||
return None
|
||||
actual = pair[target_field].to_numpy(dtype=float)
|
||||
pred = pair[forecast_field].to_numpy(dtype=float)
|
||||
ss_res = float(np.sum((actual - pred) ** 2))
|
||||
ss_tot = float(np.sum((actual - np.mean(actual)) ** 2))
|
||||
r2 = None if ss_tot <= 0 else 1.0 - (ss_res / ss_tot)
|
||||
mae = float(np.mean(np.abs(actual - pred)))
|
||||
rmse = float(np.sqrt(np.mean((actual - pred) ** 2)))
|
||||
return {"r2": r2, "mae": mae, "rmse": rmse, "samples": int(len(pair))}
|
||||
|
||||
|
||||
def _quality_point(aid, variant, target, metrics):
|
||||
p = (
|
||||
Point("forecast_metrics")
|
||||
.tag("anlagen_id", aid)
|
||||
.tag("forecast", f"prog_var_{variant}")
|
||||
.tag("target", target)
|
||||
.field("samples", int(metrics["samples"]))
|
||||
.field("mae", float(metrics["mae"]))
|
||||
.field("rmse", float(metrics["rmse"]))
|
||||
.time(datetime.datetime.utcnow(), WritePrecision.S)
|
||||
)
|
||||
if metrics["r2"] is not None:
|
||||
p.field("r2", float(metrics["r2"]))
|
||||
return p
|
||||
|
||||
|
||||
def write_quality_metrics(write_api, aid, active_variants):
|
||||
points = []
|
||||
|
||||
for variant, target in QUALITY_TARGETS.items():
|
||||
if variant not in active_variants:
|
||||
continue
|
||||
try:
|
||||
metrics = _forecast_quality(aid, variant, target)
|
||||
if metrics:
|
||||
points.append(_quality_point(aid, variant, target, metrics))
|
||||
r2_text = "nan" if metrics["r2"] is None else f"{metrics['r2']:.3f}"
|
||||
print(f"R2 Anlage {aid} prog_var_{variant}: {r2_text} / samples={metrics['samples']}")
|
||||
except Exception:
|
||||
print(f"R2 Anlage {aid} prog_var_{variant}: Berechnung fehlgeschlagen")
|
||||
traceback.print_exc()
|
||||
if points:
|
||||
try:
|
||||
write_api.write(bucket=INFLUX_BUCKET, org=INFLUX_ORG, record=points)
|
||||
except Exception:
|
||||
print(f"R2 Anlage {aid}: Schreiben der Qualitaetswerte fehlgeschlagen")
|
||||
traceback.print_exc()
|
||||
|
||||
|
||||
def run_forecast(only_anlagen_id=None, manual=False):
|
||||
configs = get_configs()
|
||||
if only_anlagen_id:
|
||||
configs = [c for c in configs if c.get("anlagen_id") == only_anlagen_id]
|
||||
client = InfluxDBClient(
|
||||
url=INFLUX_URL,
|
||||
token=INFLUX_TOKEN,
|
||||
org=INFLUX_ORG,
|
||||
timeout=INFLUX_TIMEOUT_MS,
|
||||
)
|
||||
write_api = client.write_api(write_options=SYNCHRONOUS)
|
||||
completed = []
|
||||
failures = []
|
||||
try:
|
||||
for config in configs:
|
||||
aid = config["anlagen_id"]
|
||||
try:
|
||||
data_obj = build_data_object(config, training=False)
|
||||
targets = []
|
||||
if any(active(config, n) for n in (1, 3, 10, 13, 21, 23)):
|
||||
targets.append("PV")
|
||||
if any(active(config, n) for n in (2, 3, 11, 13, 22, 23)):
|
||||
targets.append("Hausverbrauch")
|
||||
if float(config.get("batt_capacity_kwh", 0.0) or 0.0) > 0 and any(active(config, n) for n in (3, 13, 23)):
|
||||
targets.append("SOC")
|
||||
# Fail BEFORE model prediction, snapshot publication or V1/V4 plan writes.
|
||||
data_obj["telemetry_quality"] = require_recent_telemetry(data_obj, targets)
|
||||
soc_age = data_obj.get("current_soc_age_minutes")
|
||||
soc_age_text = "keine Messung" if soc_age is None else f"{soc_age:.1f} min"
|
||||
print(
|
||||
f"Forecast Anlage {aid}: Batterie-SOC {data_obj['current_soc_perc']:.1f}% "
|
||||
f"({data_obj['current_soc_source']}, Alter {soc_age_text})."
|
||||
)
|
||||
forecasts, model_errors = collect_predictions(
|
||||
data_obj, config,
|
||||
{1: v1, 2: v2, 10: v10, 11: v11, 21: v21, 22: v22}, active,
|
||||
)
|
||||
data_obj['forecast_model_status'] = model_errors
|
||||
for variant, detail in model_errors.items():
|
||||
print(f"Forecast Anlage {aid} prog_var_{variant}: unavailable ({detail['errorType']}); keine Nullwerte eingesetzt.")
|
||||
p_1, p_2, p_10, p_11, p_21, p_22 = (
|
||||
forecasts[n] for n in (1, 2, 10, 11, 21, 22)
|
||||
)
|
||||
|
||||
# ENELIX_V4_SHADOW_BRIDGE
|
||||
_v4_publish_forecasts(config, [(3,p_1,p_2),(13,p_10,p_11),(23,p_21,p_22)])
|
||||
|
||||
p_3 = v3.predict(data_obj, p_1, p_2) if active(config, 3) and p_1 and p_2 else {}
|
||||
p_13 = v13.predict(data_obj, p_10, p_11) if active(config, 13) and p_10 and p_11 else {}
|
||||
p_23 = v23.predict(data_obj, p_21, p_22) if active(config, 23) and p_21 and p_22 else {}
|
||||
|
||||
forecast_sets = [
|
||||
(1, p_1), (2, p_2), (3, p_3),
|
||||
(10, p_10), (11, p_11), (13, p_13),
|
||||
(21, p_21), (22, p_22), (23, p_23),
|
||||
]
|
||||
run_hour = str(int(datetime.datetime.now(LOCAL_TZ).strftime("%H")))
|
||||
points = []
|
||||
for variant_id, pv_src, load_src, grid_src in [
|
||||
(3, p_1, p_2, p_3),
|
||||
(13, p_10, p_11, p_13),
|
||||
(23, p_21, p_22, p_23),
|
||||
]:
|
||||
if grid_src and pv_src and load_src:
|
||||
points.extend(battery_soc_points(data_obj, variant_id, pv_src, load_src, grid_src))
|
||||
active_variants = set()
|
||||
battery_plans = data_obj.get("battery_plans", {})
|
||||
for t in data_obj["df_fut"].index:
|
||||
for n, values in forecast_sets:
|
||||
if t in values:
|
||||
field = f"prog_var_{n}"
|
||||
value = float(values[t])
|
||||
points.append(_forecast_point(aid, field, t, value))
|
||||
points.append(_snapshot_point(aid, field, run_hour, t, value))
|
||||
if n in battery_plans and t in battery_plans[n].get("battery", {}):
|
||||
points.append(_forecast_point(
|
||||
aid,
|
||||
f"prog_var_{n}_battery",
|
||||
t,
|
||||
float(battery_plans[n]["battery"][t]),
|
||||
))
|
||||
active_variants.add(n)
|
||||
if points:
|
||||
write_api.write(bucket=INFLUX_BUCKET, org=INFLUX_ORG, record=points)
|
||||
print(f"Forecast Anlage {aid}: {len(points)} Punkte geschrieben inkl. run_{run_hour}.")
|
||||
else:
|
||||
print(f"Forecast Anlage {aid}: keine aktiven Prognosen oder keine Daten.")
|
||||
write_quality_metrics(write_api, aid, active_variants)
|
||||
completed.append(aid)
|
||||
except Exception as exc:
|
||||
failures.append({"anlagen_id": aid, "error": f"{type(exc).__name__}: {exc}"})
|
||||
print(f"Forecast Anlage {aid}: Lauf fehlgeschlagen.")
|
||||
traceback.print_exc()
|
||||
finally:
|
||||
client.close()
|
||||
if failures:
|
||||
raise RuntimeError(f"Forecast-Fehler: {failures}")
|
||||
return {"completed": completed, "manual": bool(manual)}
|
||||
|
||||
|
||||
def _write_json_atomic(path, payload):
|
||||
directory = os.path.dirname(path) or "."
|
||||
os.makedirs(directory, exist_ok=True)
|
||||
tmp_path = f"{path}.tmp.{os.getpid()}"
|
||||
try:
|
||||
with open(tmp_path, "w", encoding="utf-8") as handle:
|
||||
json.dump(payload, handle)
|
||||
os.replace(tmp_path, path)
|
||||
finally:
|
||||
if os.path.exists(tmp_path):
|
||||
os.remove(tmp_path)
|
||||
|
||||
|
||||
def _read_json(path):
|
||||
try:
|
||||
with open(path, "r", encoding="utf-8") as handle:
|
||||
return json.load(handle)
|
||||
except Exception:
|
||||
return {}
|
||||
|
||||
|
||||
def _last_success_age_seconds():
|
||||
stamp = _read_json(LAST_SUCCESS_PATH).get("timestamp")
|
||||
if not stamp:
|
||||
return None
|
||||
try:
|
||||
then = datetime.datetime.fromisoformat(str(stamp).replace("Z", "+00:00"))
|
||||
now = datetime.datetime.now(datetime.timezone.utc)
|
||||
return max(0.0, (now - then).total_seconds())
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def _child_environment(training=False):
|
||||
env = os.environ.copy()
|
||||
env["PYTHONUNBUFFERED"] = "1"
|
||||
if training:
|
||||
max_threads = str(max(1, int(env.get("FORECAST_TRAINING_MAX_THREADS", "1"))))
|
||||
for name in (
|
||||
"OMP_NUM_THREADS",
|
||||
"OPENBLAS_NUM_THREADS",
|
||||
"MKL_NUM_THREADS",
|
||||
"NUMEXPR_NUM_THREADS",
|
||||
"LOKY_MAX_CPU_COUNT",
|
||||
):
|
||||
env[name] = max_threads
|
||||
return env
|
||||
|
||||
|
||||
def _run_child(mode, timeout_seconds, only_anlagen_id=None):
|
||||
command = [sys.executable, "-u", os.path.abspath(__file__), mode]
|
||||
if only_anlagen_id:
|
||||
command.append(str(only_anlagen_id))
|
||||
label = "Training" if mode == "--train-once" else "Forecast"
|
||||
print(f"{label}-Kindprozess startet (Timeout {timeout_seconds}s).")
|
||||
process = subprocess.Popen(command, env=_child_environment(training=mode == "--train-once"))
|
||||
try:
|
||||
return_code = process.wait(timeout=timeout_seconds)
|
||||
except subprocess.TimeoutExpired:
|
||||
print(f"{label}-Kindprozess hat das Zeitlimit erreicht und wird beendet.")
|
||||
process.terminate()
|
||||
try:
|
||||
process.wait(timeout=15)
|
||||
except subprocess.TimeoutExpired:
|
||||
process.kill()
|
||||
process.wait(timeout=15)
|
||||
return 124
|
||||
if return_code != 0:
|
||||
print(f"{label}-Kindprozess beendet mit Status {return_code}.")
|
||||
return return_code
|
||||
|
||||
|
||||
def run_forecast_isolated(only_anlagen_id=None, source="scheduler"):
|
||||
if not _FORECAST_PROCESS_LOCK.acquire(blocking=False):
|
||||
print(f"Forecast-Aufruf ({source}) uebersprungen: bereits ein Lauf aktiv.")
|
||||
return {"status": "busy", "source": source}
|
||||
try:
|
||||
return_code = _run_child("--forecast-once", FORECAST_TIMEOUT_SECONDS, only_anlagen_id)
|
||||
return {
|
||||
"status": "ok" if return_code == 0 else "error",
|
||||
"source": source,
|
||||
"return_code": return_code,
|
||||
}
|
||||
finally:
|
||||
_FORECAST_PROCESS_LOCK.release()
|
||||
|
||||
|
||||
def _training_worker(day_text):
|
||||
global _TRAINING_THREAD
|
||||
started = datetime.datetime.now(datetime.timezone.utc).isoformat()
|
||||
_write_json_atomic(LAST_TRAINING_PATH, {
|
||||
"day": day_text,
|
||||
"status": "running",
|
||||
"started_at": started,
|
||||
})
|
||||
return_code = 1
|
||||
error = None
|
||||
try:
|
||||
return_code = _run_child("--train-once", TRAINING_TIMEOUT_SECONDS)
|
||||
except Exception as exc:
|
||||
error = f"{type(exc).__name__}: {exc}"
|
||||
traceback.print_exc()
|
||||
finally:
|
||||
status = "ok" if return_code == 0 else ("timeout" if return_code == 124 else "error")
|
||||
payload = {
|
||||
"day": day_text,
|
||||
"status": status,
|
||||
"started_at": started,
|
||||
"finished_at": datetime.datetime.now(datetime.timezone.utc).isoformat(),
|
||||
"return_code": return_code,
|
||||
}
|
||||
if error:
|
||||
payload["error"] = error
|
||||
_write_json_atomic(LAST_TRAINING_PATH, payload)
|
||||
print(f"Nachttraining beendet: {status}.")
|
||||
_TRAINING_THREAD = None
|
||||
|
||||
|
||||
def start_training_async(day):
|
||||
global _TRAINING_THREAD
|
||||
if _TRAINING_THREAD is not None and _TRAINING_THREAD.is_alive():
|
||||
print("Nachttraining laeuft bereits.")
|
||||
return False
|
||||
day_text = day.isoformat()
|
||||
_TRAINING_THREAD = threading.Thread(target=_training_worker, args=(day_text,), daemon=True)
|
||||
_TRAINING_THREAD.start()
|
||||
return True
|
||||
|
||||
|
||||
def _watchdog_loop():
|
||||
while True:
|
||||
time.sleep(max(60, WATCHDOG_INTERVAL_SECONDS))
|
||||
try:
|
||||
age = _last_success_age_seconds()
|
||||
if age is None or age > FORECAST_STALE_SECONDS:
|
||||
age_text = "unbekannt" if age is None else f"{age / 60.0:.1f} Minuten"
|
||||
print(f"Forecast-Watchdog: letzter erfolgreicher Lauf {age_text}; neuer Lauf wird gestartet.")
|
||||
run_forecast_isolated(source="watchdog")
|
||||
except Exception:
|
||||
print("Forecast-Watchdog: Pruefung fehlgeschlagen.")
|
||||
traceback.print_exc()
|
||||
|
||||
|
||||
# BEGIN EMS RESIMULATE HTTP SERVER
|
||||
_RESIMULATE_SERVER_STARTED = False
|
||||
|
||||
|
||||
def _start_resimulate_server():
|
||||
global _RESIMULATE_SERVER_STARTED
|
||||
if _RESIMULATE_SERVER_STARTED:
|
||||
return
|
||||
_RESIMULATE_SERVER_STARTED = True
|
||||
|
||||
import json as _json
|
||||
import os as _os
|
||||
import threading as _threading
|
||||
import traceback as _traceback
|
||||
import urllib.parse as _urlparse
|
||||
from http.server import BaseHTTPRequestHandler as _BaseHTTPRequestHandler, ThreadingHTTPServer as _ThreadingHTTPServer
|
||||
|
||||
class _Handler(_BaseHTTPRequestHandler):
|
||||
def log_message(self, fmt, *args):
|
||||
return
|
||||
|
||||
def _send(self, code, body):
|
||||
raw = _json.dumps(body).encode("utf-8")
|
||||
self.send_response(code)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.send_header("Content-Length", str(len(raw)))
|
||||
self.end_headers()
|
||||
self.wfile.write(raw)
|
||||
|
||||
def do_GET(self):
|
||||
self._handle()
|
||||
|
||||
def do_POST(self):
|
||||
self._handle()
|
||||
|
||||
def _handle(self):
|
||||
try:
|
||||
parsed = _urlparse.urlparse(self.path)
|
||||
if parsed.path == "/health":
|
||||
age = _last_success_age_seconds()
|
||||
healthy = age is not None and age <= FORECAST_STALE_SECONDS
|
||||
self._send(200 if healthy else 503, {
|
||||
"status": "ok" if healthy else "stale",
|
||||
"last_success_age_seconds": age,
|
||||
"training": _read_json(LAST_TRAINING_PATH),
|
||||
})
|
||||
return
|
||||
if parsed.path not in ("/run_now", "/resimulate"):
|
||||
self._send(404, {"detail": "not found"})
|
||||
return
|
||||
params = _urlparse.parse_qs(parsed.query)
|
||||
anlagen_id = (params.get("anlagen_id") or [None])[0]
|
||||
result = run_forecast_isolated(only_anlagen_id=anlagen_id, source="http")
|
||||
result["anlagen_id"] = anlagen_id
|
||||
code = 200 if result["status"] == "ok" else (409 if result["status"] == "busy" else 500)
|
||||
self._send(code, result)
|
||||
except Exception as exc:
|
||||
_traceback.print_exc()
|
||||
self._send(500, {"detail": str(exc)})
|
||||
|
||||
port = int(_os.getenv("FORECAST_ENGINE_HTTP_PORT", "9000"))
|
||||
server = _ThreadingHTTPServer(("0.0.0.0", port), _Handler)
|
||||
thread = _threading.Thread(target=server.serve_forever, daemon=True)
|
||||
thread.start()
|
||||
print(f"Forecast Resimulate HTTP Server startet auf Port {port}.")
|
||||
# END EMS RESIMULATE HTTP SERVER
|
||||
|
||||
def main():
|
||||
_start_resimulate_server()
|
||||
global last_trained_day
|
||||
print("Forecast Engine startet.")
|
||||
threading.Thread(target=_watchdog_loop, daemon=True).start()
|
||||
run_forecast_isolated(source="startup")
|
||||
while True:
|
||||
try:
|
||||
now_local = datetime.datetime.now(LOCAL_TZ)
|
||||
if now_local.hour == TRAINING_HOUR and last_trained_day != now_local.date():
|
||||
if start_training_async(now_local.date()):
|
||||
last_trained_day = now_local.date()
|
||||
next_run = (now_local + datetime.timedelta(hours=1)).replace(minute=0, second=0, microsecond=0)
|
||||
time.sleep(max(60.0, (next_run - now_local).total_seconds()))
|
||||
run_forecast_isolated(source="scheduler")
|
||||
except Exception:
|
||||
traceback.print_exc()
|
||||
time.sleep(60)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
if len(sys.argv) >= 2 and sys.argv[1] == "--forecast-once":
|
||||
selected_anlage = sys.argv[2] if len(sys.argv) >= 3 else None
|
||||
run_forecast(only_anlagen_id=selected_anlage, manual=True)
|
||||
_write_json_atomic(LAST_SUCCESS_PATH, {
|
||||
"timestamp": datetime.datetime.now(datetime.timezone.utc).isoformat(),
|
||||
"anlagen_id": selected_anlage,
|
||||
"pid": os.getpid(),
|
||||
})
|
||||
elif len(sys.argv) >= 2 and sys.argv[1] == "--train-once":
|
||||
run_training()
|
||||
else:
|
||||
main()
|
||||
Reference in New Issue
Block a user