Об этом проекте
E-Commerce Clickstream Analytics Platform — это full-stack проект в области data engineering, обрабатывающий примерно 285 миллионов событий clickstream из Kaggle-датасета за октябрь–ноябрь 2019 года.
## Обзор
Платформа загружает сырые CSV-данные (~14 ГБ) в S3-совместимое объектное хранилище (MinIO), выполняет пакетную аналитику через Apache Spark, осуществляет потоковую обработку через Kafka и Spark Structured Streaming, запускает Hadoop MapReduce-задачи и передаёт результаты через FastAPI-бэкенд на React-дашборд.
## Архитектурные компоненты
**Приём и хранение:**
- MinIO (S3-совместимое объектное хранилище) хранит сырые CSV и выходной Parquet
- PostgreSQL 16.3 хранит агрегированные KPI и готовые к запросам данные
- Redis 7.4.1 обеспечивает слой кэширования
- Apache Kafka 3.9 (режим KRaft) служит шиной сообщений для потоковых данных
**Обработка:**
- Apache Spark 3.5 обрабатывает пакетную ETL с помощью пяти PySpark-задач: расчёт дневных KPI, анализ воронки, выявление отказов от корзины, оценка аффинности продуктов и подготовка данных для MapReduce
- Apache Hadoop 3.4.1 на YARN Streaming выполняет две MapReduce-задачи (подсчёт событий по категориям и агрегация выручки по брендам)
- Spark Structured Streaming потребляет данные из Kafka, агрегируя активные сессии каждые 5 секунд и выполняя upsert в PostgreSQL
**Оркестрация и API:**
- Apache Airflow 2.9.3 оркестрирует весь пакетный конвейер как DAG
- FastAPI (с asyncpg) предоставляет REST-эндпоинты для дашборда
- Terraform-шаблоны покрывают деплой на AWS и GCP
**Фронтенд:**
- React 18.3 с TypeScript, Tailwind CSS и Recharts отображает четыре страницы:
- Overview: карточки KPI и дневные трендовые графики
- Funnel Analysis: поток просмотра→корзина→покупка и отказы по категориям
- Product Analytics: топ продуктов, брендов и категорий
- Live Monitor: отслеживание сессий в реальном времени через Server-Sent Events
## Конвейер данных
1. **Пакетная аналитика**: сырой CSV читается из MinIO; вычисляются дневные метрики (события, пользователи, выручка, конверсия), топовые продукты/бренды/категории и записываются в PostgreSQL. Отдельная задача вычисляет процентные показатели воронки по категориям, выявляет отказы от корзины и рассчитывает product affinity lift scores (выборка 10% для производительности).
2. **MapReduce**: подготовленные Parquet-файлы конвертируются в CSV и загружаются в HDFS, где Hadoop Streaming-задачи считают события по категориям/event_type и суммируют выручку по брендам.
3. **Потоковая обработка**: Kafka producer воспроизводит CSV с настраиваемой скоростью (по умолчанию 200 событий/сек). Spark Streaming consumer агрегирует активные сессии в микробатчах по 5 секунд и сохраняет их в таблицу live_sessions.
## Развёртывание
Все 13 сервисов работают через Docker Compose. Требования: Docker Desktop >= 4.30, не менее 16 ГБ ОЗУ и 30 ГБ места на диске. Для Windows требуется WSL2 с не менее чем 12 ГБ памяти. Makefile содержит цели для запуска сервисов, загрузки данных, выполнения конвейеров и запуска тестов.
URL-адреса после запуска: Dashboard (порт 3000), FastAPI docs (порт 8000), Airflow (порт 8080), MinIO Console (порт 9001), Spark Master UI (порт 8081), HDFS NameNode (порт 9870), YARN (порт 8088).
## Структура репозитория
Ключевые директории: spark/jobs/ — пять PySpark пакетных скриптов, spark/streaming/ — Kafka consumer, kafka/producer/ — replay producer, api/ — FastAPI-сервис с маршрутизированными эндпоинтами, frontend/ — React-приложение, hadoop/mapreduce/ — скрипты Streaming-задач, airflow/dags/ — оркестрация, terraform/ — IaC-шаблоны, docker/ — определения сервисов Hadoop и Kafka.
Лицензия: MIT
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.