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.
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:
- Recolectores Autónomos: Descargar, validar y procesar los datasets de forma idempotente y tolerante a fallos.
- Cero Downtime en Actualizaciones: Actualizar la base de datos mensualmente sin interrumpir la API de consultas.
- Ingestión con Bajo Consumo de Disco: Procesar decenas de gigabytes sin necesidad de descomprimir CSVs gigantescos en el almacenamiento del servidor.
- 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
NullFilteringReaderen 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-06se 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:
PENDING ➔ DISCOVERING ➔ PROBING ➔ DOWNLOADING ➔ DOWNLOADED ➔ VALIDATING_FILE ➔ PARSING ➔ LOADING_STAGING ➔ VALIDATING_DATA ➔ PUBLISHING ➔ SUCCESS
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:
- Receita Federal (CNPJ Full): Empresas, Establecimientos, Socios (QSA), Simples Nacional/MEI, CNAEs, Municipios y Naturalezas Jurídicas.
- PGFN (Deuda Activa de la Unión): Deudores no previsionales, previsionales y FGTS con importes consolidados y estado de litigio.
- CNO (Registro Nacional de Obras): Obras civiles registradas en Brasil, superficies construidas, responsables y CNAEs vinculados.
- CAFIR (Inmuebles Rurales): Inmuebles registrados, superficie total y situación catastral.
- 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!
