ADR-0004 — Orquestación de flujos batch: Apache Airflow¶
Contexto¶
Además del flujo en tiempo real, el sistema tiene numerosos trabajos programados que necesitan orquestación robusta: limpieza horaria de lecturas, reconstrucción diaria de agregados, comprobaciones de salud de los sensores, reprocesamiento histórico bajo demanda, reentrenamiento semanal de los modelos de anomalías y de predicción, predicciones nocturnas por lotes y comprobación periódica de deriva.
Estos trabajos exigen: programación estilo cron declarativa, reintentos con espera creciente, dependencias entre tareas, capacidad de reprocesar rangos históricos con un solo comando, una interfaz para monitorizar ejecuciones, idempotencia (reejecutar no duplica datos) y control fino de concurrencia para no saturar la base de datos con muchos entrenamientos a la vez.
Decisión¶
Se adopta Apache Airflow como orquestador de todo el trabajo batch. Los flujos se definen en Python con un ejecutor distribuido que permite concurrencia real y escalado horizontal de trabajadores. Los metadatos del orquestador se guardan en la misma instancia de base de datos que el resto del sistema (en un esquema aparte), reutilizando también la infraestructura de cola ya presente. Se definen grupos de concurrencia para proteger los recursos compartidos (por ejemplo, limitar cuántos entrenamientos o cuántas escrituras concurren). Se autohospeda por completo.
Consecuencias¶
Positivas:
- Orquestación visual: la interfaz muestra el sistema completo de un vistazo, con el estado de cada ejecución, sus logs y su duración.
- Reprocesamiento histórico probado y de primera clase, imprescindible para casos como "reaplica una nueva regla de limpieza a los últimos tres meses".
- Reintentos, alertas y control de acuerdos de nivel de servicio integrados.
- Control fino de concurrencia mediante grupos de recursos, para que los distintos trabajos no se saturen entre sí.
- Ecosistema enorme de integraciones ya hechas, de modo que cualquier conexión futura es un paquete más.
Negativas o costes aceptados:
- Varios servicios adicionales (planificador, interfaz, trabajadores) con su huella de memoria.
- Curva de aprendizaje real en sus conceptos propios (recuperación de ejecuciones pasadas, reglas de disparo, mapeo dinámico de tareas).
- Las definiciones de flujo se reinterpretan con frecuencia, lo que penaliza cargar dependencias pesadas en el nivel superior del archivo.
- El paso de datos entre tareas está limitado en tamaño; los volúmenes grandes se intercambian por almacenamiento externo.
Alternativas descartadas¶
- Prefect: API moderna y muy cómoda para el desarrollador, pero con un ecosistema y una adopción menores, un reprocesamiento histórico menos pulido y su mejor experiencia ligada a la versión de pago en la nube.
- Dagster: excelente modelo basado en activos y linaje, muy adecuado para pipelines de ML, pero con un modelo mental distinto del estándar "grafo de tareas" y una adopción todavía minoritaria.
- Luigi: más simple, pero sin planificador propio ni reprocesamiento histórico real, y en clara decadencia.
- Cron más scripts: cero orquestación real: sin reintentos, sin dependencias entre trabajos, sin interfaz y sin control de concurrencia. Inviable en cuanto crecen los trabajos y sus dependencias cruzadas.