batching.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628
  1. # -*- coding: utf-8 -*-
  2. """
  3. Event-Batching für Effizienz.
  4. Sammelt Events und verarbeitet sie in Batches.
  5. """
  6. from __future__ import annotations
  7. import asyncio
  8. from dataclasses import dataclass, field
  9. from datetime import datetime
  10. from typing import TYPE_CHECKING, Any, Callable, Coroutine, Generic, TypeVar
  11. if TYPE_CHECKING:
  12. from trixy_core.events.event_data.base import EventData
  13. T = TypeVar("T")
  14. @dataclass
  15. class BatchConfig:
  16. """
  17. Konfiguration für Event-Batching.
  18. """
  19. batch_size: int = 100
  20. """Maximale Batch-Größe."""
  21. flush_interval_ms: float = 1000.0
  22. """Interval für automatisches Flushen in Millisekunden."""
  23. max_pending: int = 10000
  24. """Maximale Anzahl wartender Events."""
  25. auto_flush: bool = True
  26. """Automatisches Flushen aktivieren."""
  27. flush_on_shutdown: bool = True
  28. """Beim Shutdown flushen."""
  29. @dataclass
  30. class BatchResult:
  31. """
  32. Ergebnis einer Batch-Verarbeitung.
  33. """
  34. event_count: int
  35. """Anzahl der Events."""
  36. success_count: int
  37. """Erfolgreiche Verarbeitungen."""
  38. error_count: int
  39. """Fehlgeschlagene Verarbeitungen."""
  40. duration_ms: float
  41. """Dauer in Millisekunden."""
  42. errors: list[str] = field(default_factory=list)
  43. """Fehlermeldungen."""
  44. @dataclass
  45. class BatchStats:
  46. """
  47. Statistiken für Event-Batching.
  48. """
  49. total_events: int = 0
  50. """Gesamtzahl verarbeiteter Events."""
  51. total_batches: int = 0
  52. """Gesamtzahl verarbeiteter Batches."""
  53. total_errors: int = 0
  54. """Gesamtzahl Fehler."""
  55. avg_batch_size: float = 0.0
  56. """Durchschnittliche Batch-Größe."""
  57. avg_batch_duration_ms: float = 0.0
  58. """Durchschnittliche Batch-Dauer."""
  59. pending_count: int = 0
  60. """Aktuell wartende Events."""
  61. last_flush: datetime | None = None
  62. """Zeitpunkt des letzten Flush."""
  63. class EventBatcher:
  64. """
  65. Batcher für Event-Verarbeitung.
  66. Sammelt Events und verarbeitet sie in Batches
  67. für bessere Effizienz.
  68. Example:
  69. batcher = EventBatcher(
  70. handler=process_batch,
  71. config=BatchConfig(batch_size=50, flush_interval_ms=500),
  72. )
  73. # Events hinzufügen
  74. await batcher.add("click", event_data)
  75. await batcher.add("click", event_data)
  76. # Oder manuell flushen
  77. result = await batcher.flush()
  78. # Automatisches Flushen starten
  79. await batcher.start()
  80. # Beim Shutdown
  81. await batcher.stop()
  82. """
  83. def __init__(
  84. self,
  85. handler: Callable[[list[tuple[str, "EventData"]]], Coroutine[Any, Any, None]],
  86. config: BatchConfig | None = None,
  87. ) -> None:
  88. """
  89. Initialisiert den Batcher.
  90. Args:
  91. handler: Batch-Handler-Funktion.
  92. config: Konfiguration.
  93. """
  94. self._handler = handler
  95. self._config = config or BatchConfig()
  96. self._pending: list[tuple[str, "EventData"]] = []
  97. self._lock = asyncio.Lock()
  98. self._flush_task: asyncio.Task | None = None
  99. self._running = False
  100. self._stats = BatchStats()
  101. self._batch_durations: list[float] = []
  102. @property
  103. def pending_count(self) -> int:
  104. """Anzahl wartender Events."""
  105. return len(self._pending)
  106. @property
  107. def is_running(self) -> bool:
  108. """Ob der Batcher läuft."""
  109. return self._running
  110. async def add(
  111. self,
  112. event_name: str,
  113. event_data: "EventData",
  114. ) -> bool:
  115. """
  116. Fügt ein Event zum Batch hinzu.
  117. Args:
  118. event_name: Name des Events.
  119. event_data: Die Event-Daten.
  120. Returns:
  121. True wenn hinzugefügt.
  122. """
  123. async with self._lock:
  124. if len(self._pending) >= self._config.max_pending:
  125. # Queue voll
  126. return False
  127. self._pending.append((event_name, event_data))
  128. # Prüfen ob Batch-Größe erreicht
  129. if len(self._pending) >= self._config.batch_size:
  130. # Sofort flushen (ohne Lock halten)
  131. asyncio.create_task(self._flush_internal())
  132. return True
  133. async def add_many(
  134. self,
  135. events: list[tuple[str, "EventData"]],
  136. ) -> int:
  137. """
  138. Fügt mehrere Events hinzu.
  139. Args:
  140. events: Liste von (event_name, event_data) Tupeln.
  141. Returns:
  142. Anzahl hinzugefügter Events.
  143. """
  144. async with self._lock:
  145. space = self._config.max_pending - len(self._pending)
  146. to_add = events[:space]
  147. self._pending.extend(to_add)
  148. if len(self._pending) >= self._config.batch_size:
  149. asyncio.create_task(self._flush_internal())
  150. return len(to_add)
  151. async def flush(self) -> BatchResult:
  152. """
  153. Flusht alle wartenden Events.
  154. Returns:
  155. BatchResult.
  156. """
  157. return await self._flush_internal()
  158. async def _flush_internal(self) -> BatchResult:
  159. """
  160. Interne Flush-Implementierung.
  161. Returns:
  162. BatchResult.
  163. """
  164. async with self._lock:
  165. if not self._pending:
  166. return BatchResult(
  167. event_count=0,
  168. success_count=0,
  169. error_count=0,
  170. duration_ms=0.0,
  171. )
  172. batch = self._pending.copy()
  173. self._pending.clear()
  174. start_time = datetime.now()
  175. success_count = 0
  176. error_count = 0
  177. errors: list[str] = []
  178. try:
  179. await self._handler(batch)
  180. success_count = len(batch)
  181. except Exception as e:
  182. error_count = len(batch)
  183. errors.append(str(e))
  184. duration_ms = (datetime.now() - start_time).total_seconds() * 1000
  185. # Statistiken aktualisieren
  186. self._stats.total_events += len(batch)
  187. self._stats.total_batches += 1
  188. self._stats.total_errors += error_count
  189. self._stats.last_flush = datetime.now()
  190. self._batch_durations.append(duration_ms)
  191. if len(self._batch_durations) > 100:
  192. self._batch_durations = self._batch_durations[-50:]
  193. self._stats.avg_batch_size = self._stats.total_events / self._stats.total_batches
  194. self._stats.avg_batch_duration_ms = sum(self._batch_durations) / len(self._batch_durations)
  195. self._stats.pending_count = len(self._pending)
  196. return BatchResult(
  197. event_count=len(batch),
  198. success_count=success_count,
  199. error_count=error_count,
  200. duration_ms=duration_ms,
  201. errors=errors,
  202. )
  203. async def start(self) -> None:
  204. """Startet automatisches Flushen."""
  205. if self._running:
  206. return
  207. self._running = True
  208. async def flush_loop():
  209. while self._running:
  210. await asyncio.sleep(self._config.flush_interval_ms / 1000)
  211. if self._pending:
  212. await self._flush_internal()
  213. self._flush_task = asyncio.create_task(flush_loop())
  214. async def stop(self) -> BatchResult | None:
  215. """
  216. Stoppt den Batcher.
  217. Returns:
  218. BatchResult vom finalen Flush.
  219. """
  220. self._running = False
  221. if self._flush_task:
  222. self._flush_task.cancel()
  223. try:
  224. await self._flush_task
  225. except asyncio.CancelledError:
  226. pass
  227. self._flush_task = None
  228. if self._config.flush_on_shutdown and self._pending:
  229. return await self.flush()
  230. return None
  231. def get_stats(self) -> BatchStats:
  232. """
  233. Gibt Statistiken zurück.
  234. Returns:
  235. BatchStats.
  236. """
  237. self._stats.pending_count = len(self._pending)
  238. return self._stats
  239. def clear(self) -> int:
  240. """
  241. Löscht alle wartenden Events.
  242. Returns:
  243. Anzahl verworfener Events.
  244. """
  245. count = len(self._pending)
  246. self._pending.clear()
  247. return count
  248. class KeyedBatcher(Generic[T]):
  249. """
  250. Batcher mit Event-Gruppierung nach Schlüssel.
  251. Sammelt Events und gruppiert sie nach Schlüssel.
  252. Example:
  253. batcher = KeyedBatcher(
  254. key_func=lambda name, data: data.user_id,
  255. handler=process_user_events,
  256. )
  257. await batcher.add("click", event_data) # user_id=123
  258. await batcher.add("scroll", event_data) # user_id=123
  259. # Batch für User 123: [("click", data), ("scroll", data)]
  260. """
  261. def __init__(
  262. self,
  263. key_func: Callable[[str, "EventData"], T],
  264. handler: Callable[[T, list[tuple[str, "EventData"]]], Coroutine[Any, Any, None]],
  265. config: BatchConfig | None = None,
  266. ) -> None:
  267. """
  268. Initialisiert den KeyedBatcher.
  269. Args:
  270. key_func: Funktion zur Schlüssel-Extraktion.
  271. handler: Batch-Handler (key, events).
  272. config: Konfiguration.
  273. """
  274. self._key_func = key_func
  275. self._handler = handler
  276. self._config = config or BatchConfig()
  277. self._pending: dict[T, list[tuple[str, "EventData"]]] = {}
  278. self._lock = asyncio.Lock()
  279. self._flush_task: asyncio.Task | None = None
  280. self._running = False
  281. async def add(
  282. self,
  283. event_name: str,
  284. event_data: "EventData",
  285. ) -> bool:
  286. """
  287. Fügt ein Event hinzu.
  288. Args:
  289. event_name: Name des Events.
  290. event_data: Die Event-Daten.
  291. Returns:
  292. True wenn hinzugefügt.
  293. """
  294. key = self._key_func(event_name, event_data)
  295. async with self._lock:
  296. total_pending = sum(len(events) for events in self._pending.values())
  297. if total_pending >= self._config.max_pending:
  298. return False
  299. if key not in self._pending:
  300. self._pending[key] = []
  301. self._pending[key].append((event_name, event_data))
  302. # Prüfen ob Batch-Größe für diesen Key erreicht
  303. if len(self._pending[key]) >= self._config.batch_size:
  304. asyncio.create_task(self._flush_key(key))
  305. return True
  306. async def _flush_key(self, key: T) -> BatchResult:
  307. """
  308. Flusht Events für einen Schlüssel.
  309. Args:
  310. key: Der Schlüssel.
  311. Returns:
  312. BatchResult.
  313. """
  314. async with self._lock:
  315. if key not in self._pending or not self._pending[key]:
  316. return BatchResult(0, 0, 0, 0.0)
  317. batch = self._pending[key].copy()
  318. del self._pending[key]
  319. start_time = datetime.now()
  320. try:
  321. await self._handler(key, batch)
  322. success = len(batch)
  323. errors = 0
  324. except Exception:
  325. success = 0
  326. errors = len(batch)
  327. duration_ms = (datetime.now() - start_time).total_seconds() * 1000
  328. return BatchResult(len(batch), success, errors, duration_ms)
  329. async def flush(self) -> dict[T, BatchResult]:
  330. """
  331. Flusht alle wartenden Events.
  332. Returns:
  333. Dictionary von Key zu BatchResult.
  334. """
  335. results: dict[T, BatchResult] = {}
  336. async with self._lock:
  337. keys = list(self._pending.keys())
  338. for key in keys:
  339. results[key] = await self._flush_key(key)
  340. return results
  341. async def start(self) -> None:
  342. """Startet automatisches Flushen."""
  343. if self._running:
  344. return
  345. self._running = True
  346. async def flush_loop():
  347. while self._running:
  348. await asyncio.sleep(self._config.flush_interval_ms / 1000)
  349. if self._pending:
  350. await self.flush()
  351. self._flush_task = asyncio.create_task(flush_loop())
  352. async def stop(self) -> dict[T, BatchResult]:
  353. """
  354. Stoppt den Batcher.
  355. Returns:
  356. Results vom finalen Flush.
  357. """
  358. self._running = False
  359. if self._flush_task:
  360. self._flush_task.cancel()
  361. try:
  362. await self._flush_task
  363. except asyncio.CancelledError:
  364. pass
  365. if self._config.flush_on_shutdown:
  366. return await self.flush()
  367. return {}
  368. class TimedBatcher:
  369. """
  370. Batcher mit Zeit-basierter Gruppierung.
  371. Gruppiert Events nach Zeitfenstern.
  372. Example:
  373. batcher = TimedBatcher(
  374. window_ms=60000, # 1 Minute
  375. handler=process_minute_batch,
  376. )
  377. """
  378. def __init__(
  379. self,
  380. handler: Callable[[datetime, list[tuple[str, "EventData"]]], Coroutine[Any, Any, None]],
  381. window_ms: float = 60000.0,
  382. config: BatchConfig | None = None,
  383. ) -> None:
  384. """
  385. Initialisiert den TimedBatcher.
  386. Args:
  387. handler: Batch-Handler (window_start, events).
  388. window_ms: Zeitfenster in Millisekunden.
  389. config: Konfiguration.
  390. """
  391. self._handler = handler
  392. self._window_ms = window_ms
  393. self._config = config or BatchConfig()
  394. self._windows: dict[int, list[tuple[str, "EventData"]]] = {}
  395. self._lock = asyncio.Lock()
  396. self._running = False
  397. self._flush_task: asyncio.Task | None = None
  398. def _get_window_key(self, timestamp: datetime) -> int:
  399. """Berechnet den Window-Key für einen Zeitstempel."""
  400. epoch_ms = timestamp.timestamp() * 1000
  401. return int(epoch_ms // self._window_ms)
  402. def _window_key_to_datetime(self, key: int) -> datetime:
  403. """Konvertiert Window-Key zu Datetime."""
  404. epoch_ms = key * self._window_ms
  405. return datetime.fromtimestamp(epoch_ms / 1000)
  406. async def add(
  407. self,
  408. event_name: str,
  409. event_data: "EventData",
  410. ) -> bool:
  411. """
  412. Fügt ein Event hinzu.
  413. Args:
  414. event_name: Name des Events.
  415. event_data: Die Event-Daten.
  416. Returns:
  417. True wenn hinzugefügt.
  418. """
  419. window_key = self._get_window_key(event_data.timestamp)
  420. async with self._lock:
  421. if window_key not in self._windows:
  422. self._windows[window_key] = []
  423. self._windows[window_key].append((event_name, event_data))
  424. return True
  425. async def flush_complete_windows(self) -> int:
  426. """
  427. Flusht abgeschlossene Zeitfenster.
  428. Returns:
  429. Anzahl geflushter Windows.
  430. """
  431. current_window = self._get_window_key(datetime.now())
  432. flushed = 0
  433. async with self._lock:
  434. complete_windows = [
  435. key for key in self._windows.keys()
  436. if key < current_window
  437. ]
  438. for window_key in complete_windows:
  439. async with self._lock:
  440. if window_key not in self._windows:
  441. continue
  442. events = self._windows.pop(window_key)
  443. window_start = self._window_key_to_datetime(window_key)
  444. try:
  445. await self._handler(window_start, events)
  446. flushed += 1
  447. except Exception:
  448. # Bei Fehler zurück in die Queue
  449. async with self._lock:
  450. if window_key not in self._windows:
  451. self._windows[window_key] = []
  452. self._windows[window_key].extend(events)
  453. return flushed
  454. async def start(self) -> None:
  455. """Startet automatisches Flushen."""
  456. if self._running:
  457. return
  458. self._running = True
  459. async def flush_loop():
  460. while self._running:
  461. await asyncio.sleep(self._window_ms / 1000)
  462. await self.flush_complete_windows()
  463. self._flush_task = asyncio.create_task(flush_loop())
  464. async def stop(self) -> None:
  465. """Stoppt den Batcher."""
  466. self._running = False
  467. if self._flush_task:
  468. self._flush_task.cancel()
  469. try:
  470. await self._flush_task
  471. except asyncio.CancelledError:
  472. pass
  473. if self._config.flush_on_shutdown:
  474. # Alle Windows flushen
  475. async with self._lock:
  476. all_windows = list(self._windows.keys())
  477. for window_key in all_windows:
  478. async with self._lock:
  479. if window_key not in self._windows:
  480. continue
  481. events = self._windows.pop(window_key)
  482. window_start = self._window_key_to_datetime(window_key)
  483. await self._handler(window_start, events)