Saltar para o conteúdo
Construindo uma API para +217 Milhões de Dados Públicos (CNPJ, PGFN, CNO e CAFIR)
Engenharia de Dados Backend

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.

Construindo uma API para +217 Milhões de Dados Públicos (CNPJ, PGFN, CNO e CAFIR)

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:

  1. Coletores Autônomos: Baixar, validar e processar os datasets de forma idempotente e tolerante a falhas.
  2. Zero Downtime na Atualização: Atualizar a base de dados mensalmente sem indisponibilizar a API de consulta.
  3. Ingestão com Baixo Consumo de Disco: Processar dezenas de gigabytes sem precisar descompactar CSVs intermediários gigantescos no storage da VPS.
  4. 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 NullFilteringReader em 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-06 sã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:

PENDINGDISCOVERINGPROBINGDOWNLOADINGDOWNLOADEDVALIDATING_FILEPARSINGLOADING_STAGINGVALIDATING_DATAPUBLISHINGSUCCESS

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:

  1. Receita Federal (CNPJ Full): Empresas, Estabelecimentos, Sócios (QSA), Simples Nacional/MEI, CNAEs, Municípios e Naturezas Jurídicas.
  2. PGFN (Dívida Ativa da União): Devedores Não Previdenciários, Previdenciários e FGTS com valores consolidados e situação de ajuizamento.
  3. CNO (Cadastro Nacional de Obras): Obras civis registradas no Brasil, áreas construídas, vínculos de responsabilidade e CNAEs associados.
  4. CAFIR (Imóveis Rurais): Imóveis cadastrados, área total e situação cadastral.
  5. 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!

Última modificação