gmao/app_new/core/services/meter_analytics.py

192 lines
9.4 KiB
Python
Raw Normal View History

"""Calculs C3 sur les relevés bruts.
Les intervalles sont volontairement calculés à la demande : MeterReading reste
la source de vérité et un reset ou un changement de compteur ne peut pas être
absorbé silencieusement dans une soustraction.
"""
from dataclasses import dataclass, asdict
from datetime import date, datetime, timedelta, timezone
from statistics import median
from typing import Optional
from sqlalchemy.orm import joinedload
from ...extensions import db
from ..models.planning import (
Meter, MeterReading, MeterAlert, MeterAlertEvidence, MeterAlertRule,
MeterHeatingRegime, GasConversion, MeterTariff, CollegeClosure,
)
from .planning_service import PlanningService
@dataclass
class ConsumptionInterval:
meter_id: int
reading_start: MeterReading
reading_end: MeterReading
start_date: datetime
end_date: datetime
raw_delta: float
calendar_days: int
consumption_per_day: float
quality: str = "EXACT"
context: dict = None
def as_dict(self):
data = asdict(self)
data["reading_start"] = self.reading_start.id
data["reading_end"] = self.reading_end.id
data["start_date"] = self.start_date.isoformat()
data["end_date"] = self.end_date.isoformat()
return data
def _reading_datetime(reading):
value = reading.reading_date or datetime.min.replace(tzinfo=timezone.utc)
return value if value.tzinfo else value.replace(tzinfo=timezone.utc)
def _context_for(meter, start, end):
counts = {"scolaire": 0, "vacances": 0, "fermeture": 0, "permanence": 0, "autre": 0}
d = start.date()
while d < end.date():
if meter.housing_unit_id:
counts["scolaire"] += 1
else:
closure = CollegeClosure.query.filter(CollegeClosure.start_date <= d, CollegeClosure.end_date >= d).first()
if closure:
kind = (closure.closure_type or "vacances").lower()
key = "permanence" if "perman" in kind else ("fermeture" if "fermet" in kind else "vacances")
counts[key] += 1
elif PlanningService.get_working_hours(d) is None:
counts["fermeture"] += 1
else:
counts["scolaire"] += 1
d += timedelta(days=1)
regimes = {"NORMAL": 0, "REDUCED": 0, "STOP": 0}
rows = MeterHeatingRegime.query.filter(
MeterHeatingRegime.meter_id == meter.id,
MeterHeatingRegime.valid_from <= end.date(),
db.or_(MeterHeatingRegime.valid_to.is_(None), MeterHeatingRegime.valid_to >= start.date()),
).all()
d = start.date()
while d < end.date():
row = next((r for r in rows if r.valid_from <= d and (r.valid_to is None or d <= r.valid_to)), None)
if row:
regimes[row.regime] += 1
d += timedelta(days=1)
return {"calendar": counts, "heating": regimes}
def consumption_intervals(meter, start=None, end=None):
query = MeterReading.query.filter_by(meter_id=meter.id).order_by(MeterReading.reading_date, MeterReading.id).options(joinedload(MeterReading.read_by))
if start:
query = query.filter(MeterReading.reading_date >= start)
if end:
query = query.filter(MeterReading.reading_date <= end)
readings = query.all()
result = []
for previous, current in zip(readings, readings[1:]):
# A reset belongs to a new physical segment; never derive a delta over it.
if previous.is_reset or current.is_reset or current.value < previous.value:
continue
days = max((_reading_datetime(current) - _reading_datetime(previous)).days, 1)
interval = ConsumptionInterval(
meter_id=meter.id, reading_start=previous, reading_end=current,
start_date=_reading_datetime(previous), end_date=_reading_datetime(current),
raw_delta=current.value - previous.value, calendar_days=days,
consumption_per_day=(current.value - previous.value) / days,
)
interval.context = _context_for(meter, interval.start_date, interval.end_date)
result.append(interval)
return result
def aggregate_intervals(intervals, period="DAY"):
"""Agrégation simple par jour/semaine/mois sans modifier les relevés."""
groups = {}
for interval in intervals:
d = interval.end_date.date()
key = d if period == "DAY" else (d - timedelta(days=d.weekday()) if period == "WEEK" else (d.year, d.month))
groups.setdefault(key, []).append(interval)
return [{"period": key, "consumption": sum(i.raw_delta for i in rows), "calendar_days": sum(i.calendar_days for i in rows),
"consumption_per_day": sum(i.raw_delta for i in rows) / max(sum(i.calendar_days for i in rows), 1), "intervals": rows}
for key, rows in sorted(groups.items(), key=lambda x: x[0])]
def remainder_for_interval(meter, interval):
children = meter.children.all() if hasattr(meter.children, "all") else list(meter.children)
if not children:
return None
total = 0.0
quality = "EXACT"
for child in children:
rows = [i for i in consumption_intervals(child) if i.start_date <= interval.start_date and i.end_date >= interval.end_date]
if rows:
total += rows[-1].consumption_per_day * interval.calendar_days
continue
history = consumption_intervals(child)
if len(history) >= 2:
total += median(i.consumption_per_day for i in history[-3:]) * interval.calendar_days
quality = "ESTIMATED"
else:
quality = "PARTIAL"
return {"value": interval.raw_delta - total, "quality": quality}
def _matches(value, operator, threshold):
return {">": value > threshold, "<": value < threshold, ">=": value >= threshold, "<=": value <= threshold}[operator]
def evaluate_alert_rule(rule, interval, commit=True):
value = interval.raw_delta if rule.metric == "consumption_interval" else interval.consumption_per_day
if rule.context_filter:
context = interval.context or {}
if rule.context_filter in ("NORMAL", "REDUCED", "STOP"):
if not (context.get("heating", {}).get(rule.context_filter) or 0): return None
elif not (context.get("calendar", {}).get(rule.context_filter) or 0): return None
if not _matches(value, rule.operator, rule.threshold): return None
open_alert = MeterAlert.query.filter_by(meter_id=rule.meter_id, alert_type="MANUAL_THRESHOLD", metric=rule.metric, period=rule.period, status="OPEN").first()
if not open_alert:
open_alert = MeterAlert(meter_id=rule.meter_id, alert_type="MANUAL_THRESHOLD", level=rule.level, metric=rule.metric, period=rule.period, observed_value=value, period_start=interval.start_date.date(), period_end=interval.end_date.date())
db.session.add(open_alert); db.session.flush()
else:
open_alert.observed_value = value
db.session.add(MeterAlertEvidence(alert_id=open_alert.id, reading_start_id=interval.reading_start.id, reading_end_id=interval.reading_end.id, observed_value=value))
if commit: db.session.commit()
return open_alert
def statistical_anomaly(meter, interval, min_intervals=3, commit=True):
history = [i for i in consumption_intervals(meter) if i.end_date < interval.start_date and i.context == interval.context]
history = [i for i in history if not MeterAlertEvidence.query.filter_by(reading_end_id=i.reading_end.id, excluded_from_baseline=True).first()]
if len(history) < min_intervals: return None
reference = median(i.consumption_per_day for i in history)
if reference == 0 or interval.consumption_per_day <= reference: return None
deviation = (interval.consumption_per_day - reference) / reference * 100
alert = MeterAlert.query.filter_by(meter_id=meter.id, alert_type="STATISTICAL_ANOMALY", metric="consumption_per_day", status="OPEN").first()
if not alert:
alert = MeterAlert(meter_id=meter.id, alert_type="STATISTICAL_ANOMALY", level="WARNING", metric="consumption_per_day", period="DAY")
db.session.add(alert); db.session.flush()
alert.observed_value, alert.reference_value, alert.deviation_percent, alert.comparable_count = interval.consumption_per_day, reference, deviation, len(history)
db.session.add(MeterAlertEvidence(alert_id=alert.id, reading_start_id=interval.reading_start.id, reading_end_id=interval.reading_end.id, observed_value=interval.consumption_per_day))
if commit: db.session.commit()
return alert
def gas_coefficient(meter, on_date):
rows = GasConversion.query.filter_by(meter_id=meter.id).filter(GasConversion.valid_from <= on_date, db.or_(GasConversion.valid_to.is_(None), GasConversion.valid_to >= on_date)).order_by(GasConversion.valid_from.desc()).all()
manual = next((row for row in rows if row.origin == "MANUAL"), None)
return (manual or (rows[0] if rows else None)).coefficient_kwh_per_m3 if (manual or rows) else None
def estimated_cost(meter, interval):
tariff = MeterTariff.query.filter_by(meter_id=meter.id).filter(MeterTariff.valid_from <= interval.end_date.date(), db.or_(MeterTariff.valid_to.is_(None), MeterTariff.valid_to >= interval.end_date.date())).order_by(MeterTariff.valid_from.desc()).first()
if not tariff: return None
amount = interval.raw_delta
if meter.meter_type == "gaz" or meter.unit == "":
coefficient = gas_coefficient(meter, interval.end_date.date())
if coefficient is None and tariff.unit == "€/kWh": return None
amount *= coefficient or 1
return amount * tariff.unit_price