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

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