/
Flyer
/
library-api
Обзор
Документация
Войти
/
Flyer
/
library-api
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
data_pipelines/kappa.py
53 строки
2 KB
Alex
Task 4: Lambda and Kappa pipelines for book views
31 май 2026, 19:12
31 май 2026, 19:12
a88a1b6
Код
Авторство
О чём код?
import redis import time REDIS_HOST = "localhost" REDIS_PORT = 6379 STREAM_KEY = "book_views" KAPPA_STATS_PREFIX = "kappa:views:" r = redis.Redis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True) # Инициализация: если обработчик новый, прочитаем все старые события для восстановления состояния print("Kappa: recovering state from existing stream...") last_processed_id = r.get("kappa:last_id") if last_processed_id is None: # читаем с начала start_id = '0-0' else: start_id = last_processed_id # Загружаем все сообщения с start_id до текущего момента pending = r.xrange(STREAM_KEY, min=start_id, max='+') for entry_id, fields in pending: book_id = fields.get("book_id") if book_id: key = f"{KAPPA_STATS_PREFIX}{book_id}" r.incr(key) last_processed_id = entry_id if pending: r.set("kappa:last_id", last_processed_id) print(f"Recovered {len(pending)} events, last id {last_processed_id}") # Создаём группу потребителей для чтения новых сообщений try: r.xgroup_create(STREAM_KEY, "kappa_group", id="0", mkstream=True) except redis.ResponseError as e: if "BUSYGROUP" not in str(e): raise print("Kappa pipeline started. Reading new events...") while True: messages = r.xreadgroup("kappa_group", "kappa_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"{KAPPA_STATS_PREFIX}{book_id}" r.incr(key) count = r.get(key) print(f"Kappa: book {book_id} count {count}") r.xack(STREAM_KEY, "kappa_group", entry_id) r.set("kappa:last_id", entry_id) # сохраняем прогресс time.sleep(0.01)