Sobre el proyecto

Este proyecto ofrece un pipeline de streaming analítico en tiempo real completamente virtualizado con Docker, que demuestra conceptos de Spark Structured Streaming sin requerir cuentas en la nube. ## Qué hace Un producto en Python simula eventos clickstream enviados a un topic de Redpanda (compatible con Kafka), y un trabajo de PySpark Structured Streaming consume esos eventos, aplica agregaciones de ventana deslizante de 1 minuto con watermarking para eventos retrasados, y escribe los resultados en la consola o archivos Parquet. ## Arquitectura - **Productor Python**: Genera eventos JSON simulados de clickstream que contienen datos de URL y marca temporal - **Redpanda**: Broker compatible con Kafka ejecutándose en Docker, sin Zookeeper - **PySpark Structured Streaming**: Lee desde Kafka como DataFrame no acotado, analiza JSON con un esquema explícito, aplica ventanas de tiempo de evento con watermarks y escribe conteos agregados de visualizaciones por URL por ventana ## Conceptos clave demostrados - Leer un topic de Kafka como DataFrame no acotado con Spark Structured Streaming - Analizar payloads de eventos JSON con esquema bajo lectura - Ventanas de tiempo de evento (ventanas fijas de 1 minuto) con watermarking para manejar eventos que llegan tarde - Agregación micro-lote continua de visualizaciones por URL - Salida a consola para aprendizaje, con opción de sink Parquet para persistencia ## Estructura del proyecto ``` streaming-pipeline-kafka-spark/ ├── docker-compose.yml # Contenedores Redpanda + Spark ├── producer/ │ └── click_event_producer.py # productol simulado de clickstream ├── spark_app/ │ └── streaming_aggregation.py # trabajo Structured Streaming └── requirements.txt ``` ## Primeros pasos ```bash git clone https://github.com/Kornelius99/streaming-pipeline-kafka-spark.git cd streaming-pipeline-kafka-spark docker-compose up -d # Terminal 1: iniciar producción de eventos python producer/click_event_producer.py # Terminal 2: ejecutar trabajo 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 ``` Los resultados aparecen como conteos de visualizaciones por ventana de 1 minuto impresos en la consola, actualizándose conforme llegan nuevos eventos. ## Por qué Redpanda Redpanda implementa el protocolo Kafka para que el código de cliente existente de Kafka funcione sin modificaciones, pero se entrega como un solo contenedor sin dependencia de Zookeeper, manteniendo docker-compose simple para fines de aprendizaje. ## Licencia MIT