Saltar al contenido
Construyendo una API para +217 Millones de Datos Públicos (CNPJ, PGFN, CNO y CAFIR)
Ingeniería de Datos Backend

Construyendo una API para +217 Millones de Datos Públicos (CNPJ, PGFN, CNO y CAFIR)

Cómo diseñé una arquitectura de recolectores resilientes, particionamiento en PostgreSQL e ingestión en streaming para procesar y servir los mayores datasets abiertos de Brasil.

Construyendo una API para +217 Millones de Datos Públicos (CNPJ, PGFN, CNO y CAFIR)

Procesar datos públicos en Brasil es una verdadera prueba de fuego para cualquier ingeniero de software. Los portales gubernamentales publican volcados mensuales gigantescos, inestables, sin endpoints de delta y frecuentemente empaquetados en formatos legados con caracteres corruptos.

En este artículo, comparto cómo estructuré CNPJ-API: un ecosistema completo de recolección, procesamiento, almacenamiento y consulta vía API de los principales datos públicos brasileños — abarcando CNPJ de la Receita Federal, Deuda Activa de la PGFN, Registro Nacional de Obras (CNO), Registro de Inmuebles Rurales (CAFIR) y Créditos de la RFB.


1. Contexto: El Desafío de los Datos Públicos Brasileños

Los datasets abiertos de la Receita Federal y de la PGFN son sumamente valiosos, pero plantean desafíos operativos considerables:

  • Volumen Masivo: Tan solo el registro de CNPJ (Empresas, Establecimientos, Socios y Simples) supera los 217 millones de registros, generando más de 80 GB de datos relacionales indexados.
  • Ausencia de APIs de Cambio (Sin CDC/Delta): El fisco no informa cambios diarios; publica mensualmente un snapshot completo compuesto por decenas de archivos .zip.
  • Calidad de los Datos: Archivos CSV codificados en LATIN1, delimitados por punto y coma, con bytes nulos (\x00) e inconsistencias históricas que rompen parsers convencionales.
  • Inestabilidad de Descarga: Servidores con caídas frecuentes, tráfico restringido e interrupciones a mitad de descargas de varios gigabytes.

El objetivo era claro: ejecutar todo este pipeline de forma 100% automatizada, resiliente y eficiente en una VPS, sin recurrir a costosos clusters de Spark ni infraestructuras pesadas de Big Data.


2. Objetivo del Proyecto

Transformar múltiples dumps gubernamentales sin procesar en una plataforma unificada y consultable en milisegundos:

  1. Recolectores Autónomos: Descargar, validar y procesar los datasets de forma idempotente y tolerante a fallos.
  2. Cero Downtime en Actualizaciones: Actualizar la base de datos mensualmente sin interrumpir la API de consultas.
  3. Ingestión con Bajo Consumo de Disco: Procesar decenas de gigabytes sin necesidad de descomprimir CSVs gigantescos en el almacenamiento del servidor.
  4. API de Alto Rendimiento: Exponer endpoints REST con FastAPI para consultas estructuradas por CNPJ, Cuadro Societario (QSA), Deudas Tributarias, Obras e Inmuebles Rurales.

3. Arquitectura y Decisiones de Ingeniería

Para garantizar que el sistema no falle silenciosamente ni publique datos corruptos, la arquitectura se organizó en capas bien definidas:

       Fuentes Externas (RFB, PGFN, CNO, CAFIR)
           [ Recolector / Source Adapters ]
     (Discovery ➔ Probe ➔ Stream Download & Checksum)
             [ Pipeline de Streaming ]
       (Filtro de Bytes Nulos ➔ PostgreSQL COPY Directo)
           [ PostgreSQL 16 (Particionado) ]
   (Staging Temporal ➔ Validación Semántica ➔ Publicación Atómica)
              [ FastAPI + Caché Redis ]
                (Consultas < 50ms)

A. Streaming Unzip & Bulk Copy (Cero Desperdicio de Disco)

En lugar de descargar 5 GB de ZIPs, descomprimir 20 GB de CSVs en disco y luego importarlos a la base de datos (lo que exigiría más de 30 GB libres solo para temporales), implementé un pipeline de streaming en memoria:

  • El recolector abre el archivo ZIP en memoria mediante stream (archive.open()).
  • Los bloques pasan por un NullFilteringReader en Python que remueve bytes nulos (\x00) al vuelo.
  • El flujo se inyecta directamente en PostgreSQL con el comando nativo COPY FROM STDIN WITH (FORMAT csv, DELIMITER ';').

Esto logra una velocidad de ingestión de decenas de miles de filas por segundo, descartando los datos brutos de inmediato.

B. Particionamiento Mensual y Publicación Atómica

Para evitar que la API sirva datos incompletos durante una recolección de 2 horas, adopté particionamiento temporal por snapshot:

  • Los datos del mes 2026-06 se cargan en tablas particionadas dedicadas (cnpj.empresas_2026_06, cnpj.estabelecimentos_2026_06).
  • Una tabla de puntero (cnpj.snapshot_atual) señala el snapshot activo.
  • Validación Semántica: Antes de publicar, el sistema valida la coherencia en la cantidad de registros frente a meses anteriores (ej.: fallo automático ante caídas superiores al 30%).
  • Cambio Atómico: Solo cuando toda la carga y comprobaciones pasan con éxito, el puntero se actualiza al nuevo mes. Si la recolección falla a la mitad, la API sigue sirviendo el mes anterior sin afectación.

C. Máquina de Estados Estricta (collection_runs)

Cada recolección registra su ciclo de vida auditable en PostgreSQL:

PENDINGDISCOVERINGPROBINGDOWNLOADINGDOWNLOADEDVALIDATING_FILEPARSINGLOADING_STAGINGVALIDATING_DATAPUBLISHINGSUCCESS

Ante cualquier corte de red o cambio de schema en la fuente, el sistema registra el punto exacto de la falla, la URL, el intento y el mensaje de error.


4. Datasets Integrados

Actualmente, el ecosistema consolida 5 grandes fuentes públicas:

  1. Receita Federal (CNPJ Full): Empresas, Establecimientos, Socios (QSA), Simples Nacional/MEI, CNAEs, Municipios y Naturalezas Jurídicas.
  2. PGFN (Deuda Activa de la Unión): Deudores no previsionales, previsionales y FGTS con importes consolidados y estado de litigio.
  3. CNO (Registro Nacional de Obras): Obras civiles registradas en Brasil, superficies construidas, responsables y CNAEs vinculados.
  4. CAFIR (Inmuebles Rurales): Inmuebles registrados, superficie total y situación catastral.
  5. RFB (Créditos y Transacciones): Créditos tributarios activos, fraccionamientos de deuda y transacciones fiscales.

5. Resultados y Aprendizajes

  • +217 Millones de Registros Gestionados: Base completa de datos de Brasil ejecutándose de forma estable en PostgreSQL 16.
  • Respuestas en Sub-50ms: Consultas completas de CNPJ (socios, domicilio, situación fiscal y deudas activas) en menos de 50ms.
  • Eficiencia de Recursos: Todo el pipeline opera en una VPS convencional con Docker Compose y locks distribuidos en Redis.
  • Extensibilidad: Los adaptadores de fuente facilitan sumar nuevos orígenes de datos gubernamentales implementando discovery, descarga y parsing sin tocar el núcleo de la API.

6. Próximos Pasos

  • Dashboard de administración para observabilidad de cada recolector (SLA de actualización, fallos consecutivos y tiempos de ejecución).
  • Optimización para ingestión diferencial (Delta/Upsert) para reducir el consumo de almacenamiento mes a mes.
  • Webhooks para notificaciones automáticas de cambios de estado y situación fiscal de empresas monitoreadas.

¿Te gustó el proyecto o quieres conversar sobre ingeniería de datos y sistemas distribuidos? ¡Deja un comentario o ponte en contacto!

Última actualización