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

Этот проект предоставляет полностью Dockerизованный локальный конвейер real-time потоковой аналитики, демонстрирующий концепции Spark Structured Streaming без необходимости каких-либо облачных аккаунтов. ## Что он делает Python-продюсер генерирует эмулированные clickstream события, отправляемые в тему Redpanda (совместимую с Kafka), а задача PySpark Structured Streaming потребляет эти события, применяет 1-минутные tumbling-window агрегации с watermarking для опоздавших событий и выводит результаты в консоль или файлы Parquet. ## Архитектура - **Python Producer**: Генерирует эмулированные JSON-события clickstream, содержащие данные URL и timestamp - **Redpanda**: Kafka-совместимый брокер, работающий в Docker, не требует Zookeeper - **PySpark Structured Streaming**: Читает из Kafka как unbounded DataFrame, парсит JSON с явной schema, применяет event-time windowing с watermarks и записывает агрегированные подсчёты page-view на URL за каждое окно ## Ключевые продемонстрированные концепции - Чтение Kafka-темы как unbounded DataFrame через Spark Structured Streaming - Парсинг JSON-пейлоадов событий с schema-on-read - Event-time windowing (1-минутные tumbling windows) с watermarking для обработки опоздавших событий - Непрерывная микропакетная агрегация page-view на URL - Вывод в консоль для обучения, с опцией сохранения через Parquet sink ## Структура проекта ``` streaming-pipeline-kafka-spark/ ├── docker-compose.yml # Redpanda + Spark контейнеры ├── producer/ │ └── click_event_producer.py # симулятор clickstream продюсера ├── spark_app/ │ └── streaming_aggregation.py # Structured Streaming задача └── requirements.txt ``` ## Начало работы ```bash git clone https://github.com/Kornelius99/streaming-pipeline-kafka-spark.git cd streaming-pipeline-kafka-spark docker-compose up -d # Терминал 1: начать генерацию событий python producer/click_event_producer.py # Терминал 2: запустить Spark Streaming задачу docker-compose exec spark spark-submit \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1 \ /opt/spark_app/streaming_aggregation.py ``` Результаты появляются в виде 1-минутных windowed подсчётов page-view, печатаемых в консоль и обновляющихся по мере поступления новых событий. ## Почему Redpanda Redpanda использует протокол Kafka, поэтому существующий код Kafka-клиентов работает без изменений, но поставляется как единый контейнер без зависимости от Zookeeper, что сохраняет docker-compose простой для учебных целей. ## Лицензия MIT