El problema fundamental que aborda este artículo es la asignación eficiente de recursos computacionales en sistemas de procesamiento de streams a gran escala, específicamente Apache Flink. En entornos de hyperscaler como Netflix, con decenas de miles de jobs Flink que exhiben patrones de carga dinámicos, el aprovisionamiento estático de recursos resulta en un desperdicio significativo o en degradación del servicio durante picos de demanda. La autoescalabilidad es crucial para optimizar el costo y la latencia, pero su implementación efectiva para cargas de trabajo complejas y con estado presenta desafíos inherentes.

Históricamente, la ausencia de soluciones maduras llevó a Netflix a desarrollar un autoescalador interno. Sin embargo, la evolución de las cargas de trabajo hacia DAGs (Directed Acyclic Graphs) multi-operador y con estado, junto con la maduración de la oferta de la comunidad Flink, ha impulsado una reevaluación de la estrategia. La tesis central es que, a pesar de la inversión en una solución interna, la adopción y extensión de un proyecto de código abierto bien diseñado es una estrategia superior a largo plazo para manejar la complejidad y escala, siempre que se aborden las brechas de ingeniería específicas del entorno de producción.

Arquitectura del Sistema

El autoescalador interno original de Netflix operaba como un job de procesamiento de streams en Mantis, consumiendo métricas de Atlas (CPU, red, Kafka lag, input/consume rate) a nivel de clúster. Utilizaba una combinación de tiempo de recuperación derivado del lag, umbrales de utilización, historial de rendimiento y regresión de la tasa de entrada para decidir cuándo escalar. Este sistema escalaba el número total de TaskManagers como una unidad, lo que era adecuado para pipelines simples de un solo operador, pero ineficaz para DAGs complejos y con estado.

El autoescalador de Apache Flink, por contraste, razona desde el interior del job. Su idea clave es estimar la Tasa de Procesamiento Verdadera (TPR) de cada operador, calculando el throughput que podría sostener si estuviera completamente ocupado (throughput observado / fracción de tiempo ocupado). Recorre el grafo del job desde las fuentes, utilizando el TPR de cada operador, sus ratios de entrada/salida y una utilización objetivo para calcular el paralelismo necesario para cada vértice, asegurando que ningún operador se convierta en un cuello de botella. Esto permite un escalado granular a nivel de operador.

Para integrar el autoescalador OSS, Netflix lo implementó como un servicio Spring Boot orquestado por Temporal, un motor de flujo de trabajo duradero. Un workflow orquestador consulta el plano de control de Flink de Netflix para jobs con autoescalado habilitado y lanza un workflow de larga duración por cada job. Cada workflow de job extrae métricas por vértice del Flink JobManager, ejecuta el algoritmo de evaluación OSS y, si se toma una decisión de escalado, la entrega a un 'realizer' que aplica el cambio a través del plano de control de Flink. Esta arquitectura de 'workflow-por-job' aísla el radio de impacto de jobs problemáticos. Se realizaron modificaciones en el fork interno de Flink para optimizar la recolección de métricas (caching, filtrado server-side), preservar el 'forward chaining' (escalando subgrafos conectados como una unidad) y respetar los límites de los sinks (detectando backpressure de sinks asíncronos). El 'realizer' también incorpora comprobaciones de seguridad, como la verificación de espacio en disco para checkpoints y la prevención de escalado descendente durante evacuaciones de región.

Flujo de Autoescalado del Autoescalador OSS

  1. 1 Orquestador Temporal Consulta el plano de control de Flink para jobs con autoescalado habilitado (...
  2. 2 Workflow por Job (Temporal) Inicia un workflow de larga duración para cada job Flink.
  3. 3 Workflow por Job (Temporal) Extrae métricas por vértice del Flink JobManager (con optimizaciones de cachi...
  4. 4 Workflow por Job (Temporal) Ejecuta el algoritmo de evaluación OSS (calcula TPR, paralelismo por vértice).
  5. 5 Realizer Recibe la decisión de escalado, ejecuta comprobaciones de seguridad (ej. espa...
  6. 6 Plano de Control de Flink Actúa el cambio de escalado (savepoint, detiene, reinicia con nuevo tamaño).
CapaTecnologíaJustificación
data-processing Apache Flink Plataforma principal de procesamiento de streams para cargas de trabajo de baja latencia y alta disponibilidad.
orchestration Temporal Motor de flujo de trabajo duradero para orquestar los workflows de autoescalado por job, proporcionando aislamiento y reintentos automáticos.
observability Atlas Plataforma de telemetría interna de Netflix, fuente de métricas para el autoescalador original.
compute Spring Boot Framework para la implementación del servicio que aloja la lógica del autoescalador OSS y su integración con Temporal.
networking Kafka Sistema de mensajería para la ingesta y distribución de datos en los jobs de Flink.

Trade-offs

Ganancias
  • Eficiencia de recursos
  • Soporte para DAGs complejos y stateful
  • Reducción de costos operativos
  • Estabilidad del sistema (menos rescales)
Costes
  • Eficiencia marginal (por menor utilización objetivo)
  • Complejidad de integración (adaptación de OSS a plataforma interna)

Fundamentos Teóricos

El problema de la asignación dinámica de recursos en sistemas distribuidos de procesamiento de streams se relaciona con principios fundamentales de la teoría de colas y la optimización de sistemas. La estimación de la Tasa de Procesamiento Verdadera (TPR) de cada operador en el autoescalador OSS es una aplicación práctica de conceptos de rendimiento y capacidad de sistemas, donde se busca inferir la capacidad máxima de un componente a partir de su utilización observada y su throughput. Esto se alinea con modelos de rendimiento como los descritos por papers clásicos en teoría de colas y análisis de rendimiento de sistemas, que buscan predecir el comportamiento de un sistema bajo diferentes cargas y configuraciones.

La necesidad de un escalado granular a nivel de operador en DAGs complejos de Flink, en contraste con el escalado de clústeres monolíticos, refleja la evolución de los sistemas distribuidos hacia arquitecturas más desagregadas y conscientes de la topología de la aplicación. Esto se conecta con la investigación en scheduling y resource management en entornos de computación distribuida, donde la optimización no solo se centra en la cantidad total de recursos, sino también en su distribución y colocación para minimizar cuellos de botella y maximizar el throughput. Aunque no se cita un paper específico, los principios de control de retroalimentación (feedback control) son inherentes a cualquier sistema de autoescalado, donde las métricas del sistema actúan como entrada para un controlador que ajusta los recursos, un concepto bien establecido en la ingeniería de control.