`AlertRepository.create_many` envoyait un `INSERT` d'un seul tenant. À douze colonnes par alerte, le plafond asyncpg de 32 767 paramètres tombe à 2 730 lignes : une détection sur une fenêtre chargée échouait en `InterfaceError`, et le DAG `alertes` avec elle. Reprend le patron déjà en place dans `RecommendationRepository.create_missing`.
65 lines
2.7 KiB
Python
65 lines
2.7 KiB
Python
from collections.abc import Sequence
|
|
|
|
from sqlalchemy import select
|
|
from sqlalchemy.dialects.postgresql import insert
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from app.models.energy import Alert
|
|
|
|
# Douze colonnes par alerte, contre quatre pour une recommandation : le plafond asyncpg de
|
|
# 32 767 parametres tombe a 2 730 lignes, d'ou un lot plus petit que `recommendation.py`.
|
|
TAILLE_DE_LOT = 1000
|
|
|
|
|
|
class AlertRepository:
|
|
def __init__(self, session: AsyncSession) -> None:
|
|
self._session = session
|
|
|
|
async def list_all(
|
|
self, *, site_id: str | None = None, severity: str | None = None
|
|
) -> Sequence[Alert]:
|
|
requete = select(Alert).order_by(Alert.timestamp.desc(), Alert.alert_id.desc())
|
|
if site_id is not None:
|
|
requete = requete.where(Alert.site_id == site_id)
|
|
if severity is not None:
|
|
requete = requete.where(Alert.severity == severity)
|
|
return (await self._session.scalars(requete)).all()
|
|
|
|
async def create_many(self, alerts: Sequence[Alert]) -> Sequence[Alert]:
|
|
# `ON CONFLICT DO NOTHING` sur `uq_alert_source_reference` : rejouer la détection sur une
|
|
# fenêtre qui recouvre une exécution précédente ne doit pas dupliquer une alerte déjà
|
|
# enregistrée. `RETURNING` ne renvoie donc que les lignes effectivement insérées.
|
|
if not alerts:
|
|
return []
|
|
valeurs = [
|
|
{
|
|
"source_alert_id": alerte.source_alert_id,
|
|
"site_id": alerte.site_id,
|
|
"source": alerte.source,
|
|
"timestamp": alerte.timestamp,
|
|
"type": alerte.type,
|
|
"severity": alerte.severity,
|
|
"message": alerte.message,
|
|
"value": alerte.value,
|
|
"threshold": alerte.threshold,
|
|
"metric": alerte.metric,
|
|
"prediction_id": alerte.prediction_id,
|
|
"raw_data": alerte.raw_data,
|
|
}
|
|
for alerte in alerts
|
|
]
|
|
creees: list[Alert] = []
|
|
# Piège : asyncpg plafonne une requête à 32 767 paramètres. Une détection sur une fenêtre
|
|
# chargée dépasse ce seuil, et l'`INSERT` d'un seul tenant échouerait.
|
|
for debut in range(0, len(valeurs), TAILLE_DE_LOT):
|
|
requete = (
|
|
insert(Alert)
|
|
.values(valeurs[debut : debut + TAILLE_DE_LOT])
|
|
.on_conflict_do_nothing(constraint="uq_alert_source_reference")
|
|
.returning(Alert)
|
|
)
|
|
resultat = await self._session.execute(requete)
|
|
creees.extend(resultat.scalars().all())
|
|
await self._session.flush()
|
|
return creees
|