Об этом проекте
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.
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.