Una reconstrucción del pipeline de monitoreo de calidad del aire en México sobre un stack distinto — cambiando el CLI declarativo de Bruin por una combinación de dlt + dbt + Dagster, construida como proyecto capstone del DataTalks Data Engineering Zoomcamp 2026. El mismo problema — centralizar lecturas de OpenAQ de 306 estaciones en 121 ciudades mexicanas en tablas listas para analítica — resuelto con orquestación basada en assets y pruebas nativas de dbt en su lugar.
Arquitectura y Stack
- Ingesta → Pipelines de dlt 1.24: carga por lotes desde el archivo S3 de OpenAQ (CSV.gz particionados por
locationid/year/month) más la API v3 de OpenAQ para metadata de estaciones, archivando cada CSV.gz crudo en GCS antes de parsearlo hacia BigQuery; un productor/consumidor Redpanda compatible con Kafka transmite mediciones casi en tiempo real directamente a BigQuery - Almacenamiento → Google BigQuery: warehouse de 3 capas (
air_quality_raw→air_quality_staging→air_quality_marts), con un bucket de GCS que resguarda el archivo crudo - Transformación → dbt 1.11 (dbt-bigquery): el modelo de staging deduplica vía
ROW_NUMBER()y filtra valores atípicos, los modelos mart aplican conversión de unidades ppm→µg/m³ y categorización AQI de la OMS para PM2.5, todo validado por 15 pruebas de dbt (not_null,unique,accepted_values) - Orquestación → Assets definidos por software de Dagster 1.13 (
dagster-dbt,dagster-dlt) que conectan la ingesta de dlt directamente con el linaje de dbt, con particiones mensuales desde2024-01-01hasta el presente y una corrida programada el día 1 de cada mes - Infraestructura → Terraform (bucket GCS + datasets de BigQuery, compartido con el proyecto hermano de Bruin), Docker Compose (broker Redpanda + consola)
- CI/CD → GitHub Actions: las pruebas de dbt corren en cada push y la documentación de dbt (grafo de linaje, docs de modelos y pruebas) se publica automáticamente en GitHub Pages
Logros Técnicos Clave
Linaje de assets de extremo a extremo en Dagster — openaq_locations y openaq_measurements (assets de ingesta de dlt) alimentan directamente a stg_measurements, dim_stations y fct_city_daily_aqi (assets de dbt) como un único grafo de assets de Dagster, de modo que todo el pipeline batch corre y es observable desde un solo job programado en lugar de scripts conectados manualmente.
Capa mart normalizada por unidad y clasificada según la OMS — Factores de conversión ppm→µg/m³ específicos por contaminante (CO ×1145, NO2 ×1880, SO2 ×2620, O3 ×1960, NO ×1230, NOx ×1880) y umbrales AQI de la OMS para PM2.5 (Bueno/Moderado/Dañino para Grupos Sensibles/Dañino/Muy Dañino/Peligroso) se calculan en dbt y están respaldados por 15 pruebas de esquema que pasan correctamente.
Ingesta resiliente batch + streaming — Lógica de reintentos en S3 y una política de reintentos de Dagster en el lado batch, combinadas con manejo de límites de tasa de la API de OpenAQ y una cola de mensajes muertos (DLQ) en el consumidor de streaming de Redpanda, de modo que fallas transitorias en las fuentes no pierdan datos ni detengan el pipeline.
Documentación publicada vía CI — GitHub Actions ejecuta la suite completa de pruebas de dbt en cada push y publica el sitio de documentación de dbt generado (grafo de linaje interactivo, descripciones de modelos y columnas) en GitHub Pages, manteniendo la documentación sincronizada automáticamente con los modelos. Los resultados alimentan un dashboard de Power BI que da seguimiento a las tendencias de AQI por ciudad y al cumplimiento de las guías de la OMS.