/
Flyer
/
library-api
Обзор
Документация
Войти
/
Flyer
/
library-api
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
sales_analytics/streaming.py
53 строки
2 KB
Alex
Task 9: Streaming and batch sales analytics
02 июн 2026, 19:48
02 июн 2026, 19:48
febbadd
Код
Авторство
О чём код?
import redis import time from datetime import datetime, timedelta, timezone REDIS_HOST = "localhost" REDIS_PORT = 6379 STREAM_KEY = "sales" WINDOW_SECONDS = 300 # 5 минут r = redis.Redis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True) # Храним историю продаж как список кортежей (timestamp_epoch, price) sales_window = [] def prune_window(now_ts: float): """Удаляет события старше WINDOW_SECONDS.""" cutoff = now_ts - WINDOW_SECONDS while sales_window and sales_window[0][0] < cutoff: sales_window.pop(0) def update_window(timestamp_str: str, price: float): # timestamp_str в ISO формате ts = datetime.fromisoformat(timestamp_str).timestamp() sales_window.append((ts, price)) def get_running_average(): if not sales_window: return 0.0 total = sum(p[1] for p in sales_window) return round(total / len(sales_window), 2) print("Streaming processor started. Calculating 5-minute running average...") try: while True: # Читаем новые сообщения из потока (начиная с последнего обработанного) # Для простоты читаем все сообщения старше 0-0, но на практике нужно отслеживать позицию # Здесь будем использовать простой подход: читаем все сообщения за последний проход, # обновляем окно, печатаем среднее, затем ждём. messages = r.xrange(STREAM_KEY, min='-', max='+') # Очищаем окно и пересчитываем полностью (неэффективно, но наглядно) sales_window.clear() for _, fields in messages: ts_str = fields.get("timestamp") price = float(fields.get("price", 0)) if ts_str: update_window(ts_str, price) now_ts = time.time() prune_window(now_ts) avg = get_running_average() print(f"[{datetime.now().strftime('%H:%M:%S')}] Running average (last 5 min): {avg} (events in window: {len(sales_window)})") time.sleep(2) # обновление каждые 2 секунды except KeyboardInterrupt: print("Streaming processor stopped.")