Adrián López Rendón · proyectos

Pipeline de Monitoreo de Calidad del Aire en México

419 words 2 min read #Ingeniería de Datos#BigQuery#Kafka#Terraform

La contaminación del aire en ciudades mexicanas excede rutinariamente los lineamientos de la OMS, pero los datos de medición de más de 300 estaciones de monitoreo en 121 ciudades están dispersos y con formatos inconsistentes. Este pipeline centraliza, limpia y transforma esas lecturas en tablas listas para analítica, permitiendo a analistas y responsables de salud pública identificar puntos críticos de contaminación, dar seguimiento al cumplimiento de umbrales de la OMS y actuar sobre tendencias de calidad del aire casi en tiempo real.

Arquitectura y Stack

  • Ingesta → Doble vía: carga por lotes desde el archivo S3 de OpenAQ (CSV.gz históricos, más de 192K filas) + streaming casi en tiempo real vía un broker Redpanda compatible con Kafka
  • Almacenamiento → Google BigQuery: warehouse de 3 capas (air_quality_rawair_quality_stagingair_quality_marts), particionado por fecha y agrupado por parámetro contaminante
  • Transformación → Assets SQL declarativos vía Bruin CLI: deduplicación, normalización de unidades (ppm → µg/m³), filtrado de valores atípicos, clasificación de categoría AQI de la OMS, agregación por ciudad-día
  • Orquestación → Scheduler DAG de Bruin CLI con 17 validaciones de calidad de datos automatizadas e inline
  • Infraestructura → Terraform (bucket GCS + datasets de BigQuery), Docker Compose (broker Redpanda)
  • Visualización → Dashboard en Looker Studio con tarjetas KPI, mapa de burbujas, rankings de ciudades y tendencias de series de tiempo

Logros Técnicos Clave

Carga por lotes incremental e idempotente — Una ventana de retroceso configurable (LOOKBACK_DAYS=7) combinada con un patrón de eliminar-y-agregar asegura que cada corrida diaria procese solo datos recientes sin duplicados, reduciendo costos de cómputo en BigQuery frente a una recarga completa de tabla.

Arquitectura de ingesta dual (batch + streaming) — Par productor/consumidor de Kafka usando Redpanda y confluent-kafka, con micro-batching (200 registros o intervalo de descarga de 15 segundos) para balancear latencia contra costos de inserción en BigQuery.

Esquema dimensional en estrella — Tabla de hechos (fct_city_daily_aqi) con granularidad ciudad-día-parámetro, con banderas de cumplimiento OMS y etiquetas de categoría AQI precalculadas, junto con una tabla de dimensión dim_stations con puntuación de nivel de confiabilidad.

Aplicación automatizada de calidad de datos — 17 validaciones nativas de Bruin (not_null, unique, non_negative, positive, accepted_values) declaradas inline en cada asset SQL, combinadas con filtrado de valores atípicos y confirmaciones idempotentes de offset de Kafka.

Repositorio

github.com/sargent-mg/air-quality-bruin