/
Flyer
/
library-api
Обзор
Документация
Войти
/
Flyer
/
library-api
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
data_pipelines/lambda_batch.py
37 строк
1 KB
Alex
Task 4: Lambda and Kappa pipelines for book views
31 май 2026, 19:12
31 май 2026, 19:12
a88a1b6
Код
Авторство
О чём код?
import redis import time from datetime import datetime, timezone REDIS_HOST = "localhost" REDIS_PORT = 6379 STREAM_KEY = "book_views" BATCH_STATS_PREFIX = "batch:views:" BATCH_INTERVAL = 60 # секунд (для демонстрации) r = redis.Redis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True) def rebuild_batch(): # Читаем все сообщения из потока от начала до конца # В Redis Streams можно указать '0-0' как начало messages = r.xrange(STREAM_KEY, min='-', max='+') stats = {} for entry_id, fields in messages: book_id = fields.get("book_id") if book_id: stats[book_id] = stats.get(book_id, 0) + 1 # Сохраняем агрегаты в Redis with r.pipeline() as pipe: # Очищаем старые batch-ключи old_keys = r.keys(f"{BATCH_STATS_PREFIX}*") if old_keys: pipe.delete(*old_keys) for book_id, count in stats.items(): pipe.set(f"{BATCH_STATS_PREFIX}{book_id}", count) pipe.execute() print(f"[{datetime.now()}] Batch rebuild: {stats}") if __name__ == "__main__": print(f"Lambda Batch layer started, rebuilding every {BATCH_INTERVAL}s...") while True: rebuild_batch() time.sleep(BATCH_INTERVAL)