/
Flyer
/
library-api
Обзор
Документация
Войти
/
Flyer
/
library-api
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
task-7
data_pipelines/lambda_speed.py
36 строк
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 import json REDIS_HOST = "localhost" REDIS_PORT = 6379 STREAM_KEY = "book_views" SPEED_STATS_PREFIX = "speed:views:" r = redis.Redis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True) # Создаём группу потребителей (для надёжности, но здесь упростим) try: r.xgroup_create(STREAM_KEY, "speed_group", id="0", mkstream=True) except redis.ResponseError as e: if "BUSYGROUP" not in str(e): raise last_id = ">" # читать новые сообщения print("Lambda Speed layer started. Reading stream...") while True: # Читаем сообщения от группы "speed_group" messages = r.xreadgroup("speed_group", "speed_consumer", {STREAM_KEY: ">"}, count=10, block=1000) for stream, entries in messages: for entry_id, fields in entries: book_id = fields.get("book_id") if book_id: # Инкремент просмотров книги key = f"{SPEED_STATS_PREFIX}{book_id}" r.incr(key) print(f"Speed: book {book_id} count {r.get(key)} (event {entry_id})") # Подтверждаем обработку r.xack(STREAM_KEY, "speed_group", entry_id) # В реальной системе можно добавить периодический вывод, но здесь просто лог time.sleep(0.01)