Sobre o projeto
Este projeto oferece um pipeline de streaming analítico em tempo real totalmente containerizado com Docker, que demonstra conceitos do Spark Structured Streaming sem necessidade de contas em nuvem.
O que ele faz
Um produtor Python simula eventos de clickstream enviados a um tópico Redpanda (compatível com Kafka), e um trabalho PySpark Structured Streaming consome esses eventos, aplica agregações em janelas rolantes de 1 minuto com watermarking para eventos atrasados, e escreve os resultados no console ou em arquivos Parquet.
Arquitetura
- Produtor Python: gera eventos JSON simulados de clickstream contendo dados de URL e timestamp
- Redpanda: broker compatível com Kafka rodando em Docker, sem necessidade de Zookeeper
- PySpark Structured Streaming: lê do Kafka como um DataFrame ilimitado, analisa payloads JSON com schema explícito, aplica windowing baseado em tempo de evento com watermarks, e gera contagens agregadas de visualizações por URL por janela.
Conceitos-chave demonstrados
- Leitura de um tópico Kafka como DataFrame ilimitado com Spark Structured Streaming
- Parse de payloads JSON com schema-on-read
- Windowing baseado em tempo de evento (janelas rolantes de 1 minuto) com watermarking para lidar com eventos atrasados
- Agregação contínua em micro-lotes de visualizações de página por URL
- Saída para console (para aprendizado), com opção de sink Parquet para persistência
Estrutura do projeto
streaming-pipeline-kafka-spark/
docker-compose.yml # Contêineres Redpanda + Spark
producer/
click_event_producer.py # Produtor simulado de clickstream
spark_app/
streaming_aggregation.py # Trabalho Structured Streaming
requirements.txt
Como começar
Execute os comandos abaixo em terminal separado:
git clone https://github.com/Kornelius99/streaming-pipeline-kafka-spark.git
cd streaming-pipeline-kafka-spark
docker-compose up -d
No Terminal 1, inicie a produção de eventos:
python producer/click_event_producer.py
No Terminal 2, execute o job 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
Os resultados aparecerão como contagens de visualizações por janela de 1 minuto impressas no console, atualizando conforme novos eventos chegam.
Por que Redpanda
Redpanda usa o protocolo Kafka, permitindo que código cliente existente funcione sem modificações, mas é distribuído como um único contêiner sem dependência de Zookeeper, mantendo o docker-compose simples para fins educacionais.
Licença
MIT
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.