Construindo uma API para +217 Milhões de Dados Públicos (CNPJ, PGFN, CNO e CAFIR)
Como projetei uma arquitetura de coletores resilientes, particionamento no PostgreSQL e ingestão em streaming para processar e servir os maiores datasets abertos do Brasil.
Processar dados públicos no Brasil é um teste de fogo para qualquer engenheiro de software. Os portais governamentais disponibilizam dumps mensais gigantescos, instáveis, sem endpoints de delta e frequentemente empacotados em formatos legados com caracteres corrompidos.
Neste artigo, compartilho como estruturei o CNPJ-API: um ecossistema completo de coleta, tratamento, armazenamento e disponibilização em API dos principais dados públicos brasileiros — abrangendo CNPJ da Receita Federal, Dívida Ativa da PGFN, Cadastro Nacional de Obras (CNO), Cadastro de Imóveis Rurais (CAFIR) e Créditos da RFB.
1. Contexto: O Desafio dos Dados Públicos Brasileiros
Os datasets abertos da Receita Federal e da PGFN são riquíssimos, mas trazem desafios operacionais severos:
- Volume Massivo: Apenas o cadastro de CNPJ (Empresas, Estabelecimentos, Sócios e Simples) ultrapassa 217 milhões de registros, gerando mais de 80 GB de dados relacionais indexados.
- Ausência de APIs de Alteração (Sem CDC/Delta): A Receita não informa o que mudou dia a dia; ela publica mensalmente um full snapshot de dezenas de arquivos
.zip. - Qualidade dos Dados: Arquivos CSV com codificação
LATIN1, delimitados por ponto-e-vírgula, campos contendo bytes nulos (\x00) e inconsistências históricas que quebram parsers convencionais. - Instabilidade de Download: Servidores que sofrem quedas frequentes, lentidão de tráfego e interrupções no meio de downloads de arquivos com vários gigabytes.
O objetivo era claro: rodar todo esse pipeline de forma 100% automatizada, resiliente e eficiente em uma VPS, sem necessidade de clusters caros de Spark ou infraestruturas complexas de Big Data.
2. Objetivo do Projeto
Transformar múltiplos dumps governamentais brutos e desconexos em uma plataforma unificada e consultável em milissegundos:
- Coletores Autônomos: Baixar, validar e processar os datasets de forma idempotente e tolerante a falhas.
- Zero Downtime na Atualização: Atualizar a base de dados mensalmente sem indisponibilizar a API de consulta.
- Ingestão com Baixo Consumo de Disco: Processar dezenas de gigabytes sem precisar descompactar CSVs intermediários gigantescos no storage da VPS.
- API de Alta Performance: Disponibilizar endpoints REST via FastAPI com consultas estruturadas por CNPJ, Quadro Societário (QSA), Dívidas Tributárias, Obras e Imóveis Rurais.
3. Arquitetura e Decisões de Engenharia
Para garantir que o sistema não falhasse silenciosamente nem publicasse dados corrompidos, a arquitetura foi dividida em camadas bem definidas:
Fontes Externas (RFB, PGFN, CNO, CAFIR)
│
▼
[ Coletor / Source Adapters ]
(Discovery ➔ Probe ➔ Stream Download & Checksum)
│
▼
[ Streaming Pipeline ]
(Null Byte Filtering ➔ Direct PostgreSQL COPY)
│
▼
[ PostgreSQL 16 (Particionado) ]
(Staging Temporário ➔ Validação Semântica ➔ Publicação Atômica)
│
▼
[ FastAPI + Cache Redis ]
(Consultas < 50ms)A. Streaming Unzip & Bulk Copy (Zero Desperdício de Disco)
Ao invés de baixar 5 GB de ZIPs, descompactar 20 GB de CSVs no disco e depois enviar ao banco (o que exigiria mais de 30 GB livres só de arquivos temporários), criei um pipeline de streaming:
- O coletor abre o arquivo ZIP em memória via stream (
archive.open()). - Os chunks passam por um
NullFilteringReaderem Python, que remove bytes nulos (\x00) em voo. - O fluxo é injetado diretamente no PostgreSQL usando o comando nativo de alta velocidade
COPY FROM STDIN WITH (FORMAT csv, DELIMITER ';').
O resultado é uma taxa de ingestão de dezenas de milhares de linhas por segundo, com descarte imediato dos arquivos brutos após o processamento.
B. Particionamento Mensal e Publicação Atômica
Para que a API nunca sirva dados pela metade durante uma coleta de 2 horas, adotei particionamento temporal por snapshot:
- Os dados do mês
2026-06são carregados em tabelas particionadas dedicadas (cnpj.empresas_2026_06,cnpj.estabelecimentos_2026_06). - Uma tabela de ponteiro (
cnpj.snapshot_atual) aponta para o snapshot ativo. - Validação Semântica: Antes da publicação, o sistema valida se a quantidade de registros é consistente (ex.: falha automática se houver queda abrupta superior a 30% em relação ao mês anterior).
- Troca Atômica: Apenas quando toda a carga e integridade passam com sucesso, o ponteiro é atualizado para o novo mês. Se a coleta falhar na metade, a API continua intacta servindo o mês anterior.
C. Máquina de Estados Rigorosa (collection_runs)
Nenhuma coleta roda sem rastreabilidade. Cada execução passa pelos seguintes estados auditáveis no banco:
PENDING ➔ DISCOVERING ➔ PROBING ➔ DOWNLOADING ➔ DOWNLOADED ➔ VALIDATING_FILE ➔ PARSING ➔ LOADING_STAGING ➔ VALIDATING_DATA ➔ PUBLISHING ➔ SUCCESS
Se um download for interrompido ou um schema mudar na origem, o sistema registra exatamente o ponto da falha, a URL, a tentativa e o payload de erro, evitando repetições cegas.
4. Datasets Integrados
Atualmente, o ecossistema consolida 5 grandes fontes públicas:
- Receita Federal (CNPJ Full): Empresas, Estabelecimentos, Sócios (QSA), Simples Nacional/MEI, CNAEs, Municípios e Naturezas Jurídicas.
- PGFN (Dívida Ativa da União): Devedores Não Previdenciários, Previdenciários e FGTS com valores consolidados e situação de ajuizamento.
- CNO (Cadastro Nacional de Obras): Obras civis registradas no Brasil, áreas construídas, vínculos de responsabilidade e CNAEs associados.
- CAFIR (Imóveis Rurais): Imóveis cadastrados, área total e situação cadastral.
- RFB (Créditos e Transações): Créditos ativos tributários, parcelamentos de dívidas e transações fiscais individuais.
5. Resultados e Aprendizados
- +217 Milhões de Linhas Gerenciadas: Base completa de dados do Brasil rodando de forma estável no PostgreSQL 16.
- Consultas Sub-segundo: Consultas de CNPJ completo (incluindo quadro societário, endereço, situação fiscal e dívidas ativas) retornando em menos de 50ms.
- Eficiência de Recursos: Todo o pipeline roda em uma VPS convencional com Docker Compose, utilizando locks distribuídos no Redis para evitar concorrência acidental entre coletores.
- Resiliência a Mudanças: A separação em Source Adapters permite plugar novas fontes governamentais apenas implementando os passos de discovery, download e parsing sem alterar o núcleo da API.
6. Próximos Passos
Os próximos passos do projeto incluem:
- Implementação de um dashboard administrativo para monitoramento visual do status de cada coletor (SLA de atualização, falhas consecutivas e tempo de execução).
- Otimização do pipeline para ingestão diferencial (Delta/Upsert) visando reduzir ainda mais o consumo de armazenamento entre meses consecutivos.
- Exposição de webhooks para notificação automática de alterações cadastrais e situação fiscal de CNPJs monitorados.
Gostou do projeto ou quer trocar uma ideia sobre engenharia de dados e sistemas distribuídos? Deixe um comentário ou entre em contato!
