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.
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