909 lines
36 KiB
Python
909 lines
36 KiB
Python
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()
|