Об этом проекте

RealTime-ECommerce-Analytics (RTEA) — высокопроизводительный потоковый пайплайн на Python, который захватывает кликстрим-данные e-commerce из REST API и направляет их через Confluent Cloud Kafka в ClickHouse для аналитических запросов в реальном времени. ## Архитектура Пайплайн состоит из четырёх слоёв: **Источник данных** — генератор кликстрим-событий в виде API, размещённый на Railway.app, выдающий примерно 10 000 сообщений в секунду в формате JSON. **Брокер сообщений** — Confluent Cloud Kafka с топиком `clickstream-events` (6 партиций), аутентификацией SASL_SSL, сжатием Snappy и retention 7 дней. **Слой обработки** — сервис-продюсер на базе `kafka-python` с оптимизированными настройками (батчи 16 КБ, сжатие gzip, буфер 32 МБ, ограничение скорости до 10K msg/сек) и сервис-консьюзер с пакетной обработкой (1000 событий за батч, автокоммит оффсетов, трансформация полей). **Хранилище для аналитики** — ClickHouse Cloud с таблицей на движке MergeTree, упорядоченной по `(timestamp, user_id)`, сжатием LZ4 и TTL 30 дней. ## Поток данных 1. Продюсер каждую секунду опрашивает внешний API, собирает сообщения в очередь и публикует батчами в топик Kafka. 2. Несколько инстансов консьюзеров подписываются на назначенные партиции, десериализуют JSON, преобразуют временные метки и поля, и пишут батчами в ClickHouse. 3. Пользователи запрашивают данные из таблицы `ecommerce.clickstream_events` через SQL. ## Ключевые конфигурации ### Продюсер - Размер батча: 16 КБ - Сжатие: gzip - Буфер памяти: 32 МБ - Повторные попытки: 3 ### Консьюзер - Размер батча: 1000 событий - Таймаут батча: 5 секунд - Автокоммит оффсетов включён ### Таблица ClickHouse ```sql CREATE TABLE clickstream_events ( user_id String, session_id String, timestamp UInt64, event_type String, product_id String, product_category String, price Float64, quantity Int32, source String, received_time DateTime DEFAULT now() ) ENGINE = MergeTree() ORDER BY (timestamp, user_id) ``` ## Сообщаемая производительность - Устойчивая пропускная способность: 650+ событий/секунду - Всего обработано событий: 510 414+ - Сквозная латентность менее секунды ## Быстрый старт 1. Установите зависимости: `uv sync` 2. Скопируйте `.env.example` в `.env` и укажите учетные данные Confluent Cloud и ClickHouse 3. Запустите продюсер: `uv run python services/producer/src/simple_producer.py` 4. Запустите консьюзер: `uv run python services/consumer/src/simple_consumer.py` 5. Проверьте данные: `uv run python check_clickhouse.py` ## Устранение неполадок Репозиторий включает диагностические скрипты (`debug_kafka.py`, `test_clickhouse_insert.py`) и охватывает типичные сценарии отказов: таймауты подключения к API, отставание консьюмера Kafka, ошибки десериализации сообщений и несовпадения схемы/вставки в ClickHouse.