← voltar aos projetos

Engenharia de dados · CDC · Data Lake

OpenSky CDC DataLake

Captura e replicação de dados de tráfego aéreo em tempo real utilizando Change Data Capture (CDC), com ingestão via Kafka/Debezium, processamento em Python e armazenamento em camadas Bronze e Silver de um Data Lake.

Status
Concluído
Stack
PostgreSQL · Debezium · Kafka · MinIO · Python · Docker
Código

Contexto

Projeto de engenharia de dados focado em replicar, em tempo real, o estado de aeronaves em voo capturado pela OpenSky Network — posição, altitude e velocidade. O desafio central é que a API da OpenSky não expõe um changelog nativo: ela sempre retorna o estado atual das aeronaves a cada consulta. Por isso, o CDC não atua diretamente sobre a fonte, mas sobre uma camada de staging em PostgreSQL, onde as inserções e atualizações geradas pelo polling da API se tornam o gatilho real do stream de eventos.

Fluxo da arquitetura

Um script Python consulta periodicamente a API da OpenSky e grava os estados das aeronaves em uma tabela de staging no PostgreSQL. O Debezium monitora o log de transação desse banco e publica cada alteração como evento no Kafka. A partir daí o fluxo se divide em duas frentes: um S3 Sink Connector do Kafka Connect grava os eventos crus em JSON na camada Bronze do MinIO, enquanto um consumidor Python assíncrono lê os mesmos tópicos, limpa os dados (remoção de NaN, tipagem) e grava arquivos particionados em Parquet na camada Silver.

OpenSky Network API polling PostgreSQL staging estados brutos log de transação Debezium source connector detecta mudanças Apache Kafka tópicos de eventos de mudança S3 Sink Connector → MinIO Bronze (JSON) Consumidor Python → MinIO Silver (Parquet)

O que foi desenvolvido

  • Infraestrutura completa orquestrada via Docker Compose (PostgreSQL, Kafka, Zookeeper, MinIO, Kafka Connect e app Python)
  • Aplicação de ingestão em Python com polling real da API da OpenSky Network para a tabela de staging
  • Conector Source do Debezium configurado para escutar alterações na tabela de staging do PostgreSQL
  • Conector Sink S3 do Kafka Connect gravando os eventos crus na camada Bronze do MinIO
  • Consumidor Python assíncrono com Pandas/PyArrow para limpeza e conversão dos eventos em Parquet
  • Camada Silver com arquivos particionados por data, prontos para consultas analíticas
🔧 Os próximos passos do roadmap incluem adicionar Apache Iceberg + Trino para consultas SQL nativas sobre o Data Lake, e construir dashboards no Apache Superset apontando para as camadas Silver/Ouro.