Files
Enelix-EMS/services/netplan-v4/acceptance/forecast-src/main.py
T

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()