Об этом проекте
Этот проект предоставляет полностью 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
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.