| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628 |
- # -*- coding: utf-8 -*-
- """
- Event-Batching für Effizienz.
- Sammelt Events und verarbeitet sie in Batches.
- """
- from __future__ import annotations
- import asyncio
- from dataclasses import dataclass, field
- from datetime import datetime
- from typing import TYPE_CHECKING, Any, Callable, Coroutine, Generic, TypeVar
- if TYPE_CHECKING:
- from trixy_core.events.event_data.base import EventData
- T = TypeVar("T")
- @dataclass
- class BatchConfig:
- """
- Konfiguration für Event-Batching.
- """
- batch_size: int = 100
- """Maximale Batch-Größe."""
- flush_interval_ms: float = 1000.0
- """Interval für automatisches Flushen in Millisekunden."""
- max_pending: int = 10000
- """Maximale Anzahl wartender Events."""
- auto_flush: bool = True
- """Automatisches Flushen aktivieren."""
- flush_on_shutdown: bool = True
- """Beim Shutdown flushen."""
- @dataclass
- class BatchResult:
- """
- Ergebnis einer Batch-Verarbeitung.
- """
- event_count: int
- """Anzahl der Events."""
- success_count: int
- """Erfolgreiche Verarbeitungen."""
- error_count: int
- """Fehlgeschlagene Verarbeitungen."""
- duration_ms: float
- """Dauer in Millisekunden."""
- errors: list[str] = field(default_factory=list)
- """Fehlermeldungen."""
- @dataclass
- class BatchStats:
- """
- Statistiken für Event-Batching.
- """
- total_events: int = 0
- """Gesamtzahl verarbeiteter Events."""
- total_batches: int = 0
- """Gesamtzahl verarbeiteter Batches."""
- total_errors: int = 0
- """Gesamtzahl Fehler."""
- avg_batch_size: float = 0.0
- """Durchschnittliche Batch-Größe."""
- avg_batch_duration_ms: float = 0.0
- """Durchschnittliche Batch-Dauer."""
- pending_count: int = 0
- """Aktuell wartende Events."""
- last_flush: datetime | None = None
- """Zeitpunkt des letzten Flush."""
- class EventBatcher:
- """
- Batcher für Event-Verarbeitung.
- Sammelt Events und verarbeitet sie in Batches
- für bessere Effizienz.
- Example:
- batcher = EventBatcher(
- handler=process_batch,
- config=BatchConfig(batch_size=50, flush_interval_ms=500),
- )
- # Events hinzufügen
- await batcher.add("click", event_data)
- await batcher.add("click", event_data)
- # Oder manuell flushen
- result = await batcher.flush()
- # Automatisches Flushen starten
- await batcher.start()
- # Beim Shutdown
- await batcher.stop()
- """
- def __init__(
- self,
- handler: Callable[[list[tuple[str, "EventData"]]], Coroutine[Any, Any, None]],
- config: BatchConfig | None = None,
- ) -> None:
- """
- Initialisiert den Batcher.
- Args:
- handler: Batch-Handler-Funktion.
- config: Konfiguration.
- """
- self._handler = handler
- self._config = config or BatchConfig()
- self._pending: list[tuple[str, "EventData"]] = []
- self._lock = asyncio.Lock()
- self._flush_task: asyncio.Task | None = None
- self._running = False
- self._stats = BatchStats()
- self._batch_durations: list[float] = []
- @property
- def pending_count(self) -> int:
- """Anzahl wartender Events."""
- return len(self._pending)
- @property
- def is_running(self) -> bool:
- """Ob der Batcher läuft."""
- return self._running
- async def add(
- self,
- event_name: str,
- event_data: "EventData",
- ) -> bool:
- """
- Fügt ein Event zum Batch hinzu.
- Args:
- event_name: Name des Events.
- event_data: Die Event-Daten.
- Returns:
- True wenn hinzugefügt.
- """
- async with self._lock:
- if len(self._pending) >= self._config.max_pending:
- # Queue voll
- return False
- self._pending.append((event_name, event_data))
- # Prüfen ob Batch-Größe erreicht
- if len(self._pending) >= self._config.batch_size:
- # Sofort flushen (ohne Lock halten)
- asyncio.create_task(self._flush_internal())
- return True
- async def add_many(
- self,
- events: list[tuple[str, "EventData"]],
- ) -> int:
- """
- Fügt mehrere Events hinzu.
- Args:
- events: Liste von (event_name, event_data) Tupeln.
- Returns:
- Anzahl hinzugefügter Events.
- """
- async with self._lock:
- space = self._config.max_pending - len(self._pending)
- to_add = events[:space]
- self._pending.extend(to_add)
- if len(self._pending) >= self._config.batch_size:
- asyncio.create_task(self._flush_internal())
- return len(to_add)
- async def flush(self) -> BatchResult:
- """
- Flusht alle wartenden Events.
- Returns:
- BatchResult.
- """
- return await self._flush_internal()
- async def _flush_internal(self) -> BatchResult:
- """
- Interne Flush-Implementierung.
- Returns:
- BatchResult.
- """
- async with self._lock:
- if not self._pending:
- return BatchResult(
- event_count=0,
- success_count=0,
- error_count=0,
- duration_ms=0.0,
- )
- batch = self._pending.copy()
- self._pending.clear()
- start_time = datetime.now()
- success_count = 0
- error_count = 0
- errors: list[str] = []
- try:
- await self._handler(batch)
- success_count = len(batch)
- except Exception as e:
- error_count = len(batch)
- errors.append(str(e))
- duration_ms = (datetime.now() - start_time).total_seconds() * 1000
- # Statistiken aktualisieren
- self._stats.total_events += len(batch)
- self._stats.total_batches += 1
- self._stats.total_errors += error_count
- self._stats.last_flush = datetime.now()
- self._batch_durations.append(duration_ms)
- if len(self._batch_durations) > 100:
- self._batch_durations = self._batch_durations[-50:]
- self._stats.avg_batch_size = self._stats.total_events / self._stats.total_batches
- self._stats.avg_batch_duration_ms = sum(self._batch_durations) / len(self._batch_durations)
- self._stats.pending_count = len(self._pending)
- return BatchResult(
- event_count=len(batch),
- success_count=success_count,
- error_count=error_count,
- duration_ms=duration_ms,
- errors=errors,
- )
- async def start(self) -> None:
- """Startet automatisches Flushen."""
- if self._running:
- return
- self._running = True
- async def flush_loop():
- while self._running:
- await asyncio.sleep(self._config.flush_interval_ms / 1000)
- if self._pending:
- await self._flush_internal()
- self._flush_task = asyncio.create_task(flush_loop())
- async def stop(self) -> BatchResult | None:
- """
- Stoppt den Batcher.
- Returns:
- BatchResult vom finalen Flush.
- """
- self._running = False
- if self._flush_task:
- self._flush_task.cancel()
- try:
- await self._flush_task
- except asyncio.CancelledError:
- pass
- self._flush_task = None
- if self._config.flush_on_shutdown and self._pending:
- return await self.flush()
- return None
- def get_stats(self) -> BatchStats:
- """
- Gibt Statistiken zurück.
- Returns:
- BatchStats.
- """
- self._stats.pending_count = len(self._pending)
- return self._stats
- def clear(self) -> int:
- """
- Löscht alle wartenden Events.
- Returns:
- Anzahl verworfener Events.
- """
- count = len(self._pending)
- self._pending.clear()
- return count
- class KeyedBatcher(Generic[T]):
- """
- Batcher mit Event-Gruppierung nach Schlüssel.
- Sammelt Events und gruppiert sie nach Schlüssel.
- Example:
- batcher = KeyedBatcher(
- key_func=lambda name, data: data.user_id,
- handler=process_user_events,
- )
- await batcher.add("click", event_data) # user_id=123
- await batcher.add("scroll", event_data) # user_id=123
- # Batch für User 123: [("click", data), ("scroll", data)]
- """
- def __init__(
- self,
- key_func: Callable[[str, "EventData"], T],
- handler: Callable[[T, list[tuple[str, "EventData"]]], Coroutine[Any, Any, None]],
- config: BatchConfig | None = None,
- ) -> None:
- """
- Initialisiert den KeyedBatcher.
- Args:
- key_func: Funktion zur Schlüssel-Extraktion.
- handler: Batch-Handler (key, events).
- config: Konfiguration.
- """
- self._key_func = key_func
- self._handler = handler
- self._config = config or BatchConfig()
- self._pending: dict[T, list[tuple[str, "EventData"]]] = {}
- self._lock = asyncio.Lock()
- self._flush_task: asyncio.Task | None = None
- self._running = False
- async def add(
- self,
- event_name: str,
- event_data: "EventData",
- ) -> bool:
- """
- Fügt ein Event hinzu.
- Args:
- event_name: Name des Events.
- event_data: Die Event-Daten.
- Returns:
- True wenn hinzugefügt.
- """
- key = self._key_func(event_name, event_data)
- async with self._lock:
- total_pending = sum(len(events) for events in self._pending.values())
- if total_pending >= self._config.max_pending:
- return False
- if key not in self._pending:
- self._pending[key] = []
- self._pending[key].append((event_name, event_data))
- # Prüfen ob Batch-Größe für diesen Key erreicht
- if len(self._pending[key]) >= self._config.batch_size:
- asyncio.create_task(self._flush_key(key))
- return True
- async def _flush_key(self, key: T) -> BatchResult:
- """
- Flusht Events für einen Schlüssel.
- Args:
- key: Der Schlüssel.
- Returns:
- BatchResult.
- """
- async with self._lock:
- if key not in self._pending or not self._pending[key]:
- return BatchResult(0, 0, 0, 0.0)
- batch = self._pending[key].copy()
- del self._pending[key]
- start_time = datetime.now()
- try:
- await self._handler(key, batch)
- success = len(batch)
- errors = 0
- except Exception:
- success = 0
- errors = len(batch)
- duration_ms = (datetime.now() - start_time).total_seconds() * 1000
- return BatchResult(len(batch), success, errors, duration_ms)
- async def flush(self) -> dict[T, BatchResult]:
- """
- Flusht alle wartenden Events.
- Returns:
- Dictionary von Key zu BatchResult.
- """
- results: dict[T, BatchResult] = {}
- async with self._lock:
- keys = list(self._pending.keys())
- for key in keys:
- results[key] = await self._flush_key(key)
- return results
- async def start(self) -> None:
- """Startet automatisches Flushen."""
- if self._running:
- return
- self._running = True
- async def flush_loop():
- while self._running:
- await asyncio.sleep(self._config.flush_interval_ms / 1000)
- if self._pending:
- await self.flush()
- self._flush_task = asyncio.create_task(flush_loop())
- async def stop(self) -> dict[T, BatchResult]:
- """
- Stoppt den Batcher.
- Returns:
- Results vom finalen Flush.
- """
- self._running = False
- if self._flush_task:
- self._flush_task.cancel()
- try:
- await self._flush_task
- except asyncio.CancelledError:
- pass
- if self._config.flush_on_shutdown:
- return await self.flush()
- return {}
- class TimedBatcher:
- """
- Batcher mit Zeit-basierter Gruppierung.
- Gruppiert Events nach Zeitfenstern.
- Example:
- batcher = TimedBatcher(
- window_ms=60000, # 1 Minute
- handler=process_minute_batch,
- )
- """
- def __init__(
- self,
- handler: Callable[[datetime, list[tuple[str, "EventData"]]], Coroutine[Any, Any, None]],
- window_ms: float = 60000.0,
- config: BatchConfig | None = None,
- ) -> None:
- """
- Initialisiert den TimedBatcher.
- Args:
- handler: Batch-Handler (window_start, events).
- window_ms: Zeitfenster in Millisekunden.
- config: Konfiguration.
- """
- self._handler = handler
- self._window_ms = window_ms
- self._config = config or BatchConfig()
- self._windows: dict[int, list[tuple[str, "EventData"]]] = {}
- self._lock = asyncio.Lock()
- self._running = False
- self._flush_task: asyncio.Task | None = None
- def _get_window_key(self, timestamp: datetime) -> int:
- """Berechnet den Window-Key für einen Zeitstempel."""
- epoch_ms = timestamp.timestamp() * 1000
- return int(epoch_ms // self._window_ms)
- def _window_key_to_datetime(self, key: int) -> datetime:
- """Konvertiert Window-Key zu Datetime."""
- epoch_ms = key * self._window_ms
- return datetime.fromtimestamp(epoch_ms / 1000)
- async def add(
- self,
- event_name: str,
- event_data: "EventData",
- ) -> bool:
- """
- Fügt ein Event hinzu.
- Args:
- event_name: Name des Events.
- event_data: Die Event-Daten.
- Returns:
- True wenn hinzugefügt.
- """
- window_key = self._get_window_key(event_data.timestamp)
- async with self._lock:
- if window_key not in self._windows:
- self._windows[window_key] = []
- self._windows[window_key].append((event_name, event_data))
- return True
- async def flush_complete_windows(self) -> int:
- """
- Flusht abgeschlossene Zeitfenster.
- Returns:
- Anzahl geflushter Windows.
- """
- current_window = self._get_window_key(datetime.now())
- flushed = 0
- async with self._lock:
- complete_windows = [
- key for key in self._windows.keys()
- if key < current_window
- ]
- for window_key in complete_windows:
- async with self._lock:
- if window_key not in self._windows:
- continue
- events = self._windows.pop(window_key)
- window_start = self._window_key_to_datetime(window_key)
- try:
- await self._handler(window_start, events)
- flushed += 1
- except Exception:
- # Bei Fehler zurück in die Queue
- async with self._lock:
- if window_key not in self._windows:
- self._windows[window_key] = []
- self._windows[window_key].extend(events)
- return flushed
- async def start(self) -> None:
- """Startet automatisches Flushen."""
- if self._running:
- return
- self._running = True
- async def flush_loop():
- while self._running:
- await asyncio.sleep(self._window_ms / 1000)
- await self.flush_complete_windows()
- self._flush_task = asyncio.create_task(flush_loop())
- async def stop(self) -> None:
- """Stoppt den Batcher."""
- self._running = False
- if self._flush_task:
- self._flush_task.cancel()
- try:
- await self._flush_task
- except asyncio.CancelledError:
- pass
- if self._config.flush_on_shutdown:
- # Alle Windows flushen
- async with self._lock:
- all_windows = list(self._windows.keys())
- for window_key in all_windows:
- async with self._lock:
- if window_key not in self._windows:
- continue
- events = self._windows.pop(window_key)
- window_start = self._window_key_to_datetime(window_key)
- await self._handler(window_start, events)
|