Flink Autoscaler Netflix: el ahorro de USD 1,1M al año
En pocas palabras: Netflix está migrando a un autoscaler open source de Apache Flink que ajusta el paralelismo por operador en vez de por clúster completo, para cubrir sus más de 30.000 jobs de streaming en AWS. Un equipo interno logró un ahorro del 58% en cómputo, unos USD 1,1 millones al año.
Netflix está migrando el escalado automático de sus más de 30.000 jobs de streaming en Apache Flink a un autoscaler open source que decide por operador, no por clúster completo. Un equipo interno redujo el gasto anualizado en cómputo Flink un 58%, un ahorro cercano a USD 1,1 millones al año, según reportó InfoQ.
El Flink Autoscaler es un componente open source de Apache Flink que ajusta el paralelismo de cada operador de un job de streaming según su throughput real y su busy time, en vez de escalar el clúster entero como una sola unidad. Netflix lo está adoptando para reemplazar el autoscaler interno que construyó sobre Mantis desde 2019, con el objetivo de cubrir sus más de 30.000 jobs en múltiples regiones de AWS.
En este artículo:
- En 30 segundos
- ¿Qué problema tenía el autoscaler interno de Netflix hasta ahora?
- ¿Cómo funciona el nuevo Apache Flink Autoscaler basado en operadores?
- ¿Cuánto ahorra Netflix con el Flink Autoscaler?
- ¿Qué es el proyecto DS2 y por qué es la base técnica del autoscaler?
- ¿Cómo integró Netflix el autoscaler con su control plane interno?
- ¿En qué se diferencia de otros autoscalers como KEDA?
- ¿Por qué Netflix usa un target de utilización de 0,45 en vez del 0,7 por defecto?
- ¿Qué sigue para Netflix con Flink 2 y el estado desagregado?
- Qué significa esto para equipos de datos en Latinoamérica
- Errores comunes al pensar el autoscaling de streaming
- Preguntas Frecuentes
- Conclusión
- Fuentes
En 30 segundos
- Netflix opera más de 30.000 streaming jobs en Apache Flink en varias regiones de AWS.
- El nuevo Flink Autoscaler calcula el paralelismo por operador individual, no por clúster completo.
- Un equipo interno ahorró 58% en cómputo Flink, unos USD 1,1 millones al año.
- Netflix no usa el Flink Kubernetes Operator: integró el autoscaler con un servicio propio en Spring Boot y Temporal.
- La empresa fija un utilization target de 0,45, por debajo del 0,7 que recomienda la comunidad de Flink.
¿Qué problema tenía el autoscaler interno de Netflix hasta ahora?
El autoscaler original de Netflix, construido sobre Mantis alrededor de 2019, escalaba el clúster completo como una sola unidad, sin distinguir entre operadores dentro del mismo job. Tomaba telemetría a nivel de clúster (CPU, uso de red, Kafka lag, input rate y consume rate) y ajustaba la cantidad total de TaskManagers, lo que le permitió reducir el uso de recursos entre 25% y 45% en miles de pipelines.
El límite estaba en la unidad de escalado. Ponele que tenés un job con un join pesado que necesita el triple de paralelismo que el resto del pipeline: con el autoscaler viejo, todo el clúster escalaba parejo, así que terminabas sobreaprovisionando las partes livianas para cubrir la parte pesada. Para pipelines stateful con branches, joins y terabytes de estado, esa limitación se volvió cara.
¿Cómo funciona el nuevo Apache Flink Autoscaler basado en operadores?
El nuevo Apache Flink Autoscaler usa métricas que expone el propio job en ejecución para estimar la tasa de procesamiento real de cada operador, a partir de throughput y busy time. Con eso recorre el grafo del job y calcula el paralelismo necesario para cada vértice por separado, en vez de tratar el clúster como un bloque único. Esto se conecta con lo que analizamos en los pipelines de CI/CD que sostienen este tipo de infraestructura.
¿Y qué gana Netflix concretamente con este cambio de unidad de medida? Deja de pagar de más por operadores livianos solo para sostener a los pesados. El enfoque está descrito en un paper que aborda el autoscaling de jobs heterogéneos y el costo de rescalar aplicaciones stateful, la base técnica detrás de todo el Flink Autoscaler que Netflix está terminando de adoptar.
¿Cuánto ahorra Netflix con el Flink Autoscaler?
Un equipo de Netflix redujo el gasto anualizado en cómputo Flink un 58%, un ahorro de aproximadamente USD 1,1 millones al año, según el reporte publicado por InfoQ. El dato corresponde a un equipo puntual, no a un promedio de toda la flota, pero da una idea del margen que había en el modelo cluster-level.
Con más de 30.000 jobs corriendo en simultáneo en varias regiones de AWS, cualquier mejora porcentual en eficiencia de cómputo se multiplica rápido. Acá va la tabla comparando los dos enfoques:
| Aspecto | Autoscaler original (Mantis) | Flink Autoscaler (nuevo) |
|---|---|---|
| Unidad de escalado | Clúster completo | Operador individual (por vértice) |
| Métricas usadas | CPU, red, Kafka lag, input/consume rate | Throughput y busy time por operador |
| Reducción de recursos reportada | 25% a 45% | Hasta 58% en un equipo |
| Utilization target | No aplica (cluster-level) | 0,45 en Netflix (default comunitario: 0,7) |

¿Qué es el proyecto DS2 y por qué es la base técnica del autoscaler?
DS2 es el proyecto de investigación en el que se apoya la técnica de estimación de paralelismo por operador que usa el nuevo autoscaler. Un investigador de sistemas involucrado en el trabajo contó que el equipo exploró primero análisis de ruta crítica antes de adoptar DS2 como baseline más simple, y que “resultó que esta idea muy simple funcionaba realmente bien”. Esa idea terminó integrada al trabajo de autoscaling de Flink.
¿Cómo integró Netflix el autoscaler con su control plane interno?
Netflix no desplegó el autoscaler directo a través del Flink Kubernetes Operator. En cambio, construyó un servicio en Spring Boot que usa workflows de Temporal para aislar la decisión de autoscaling de cada job por separado, según describe el blog técnico de Netflix.
Subís el servicio, este consulta las métricas del JobManager, calcula el paralelismo objetivo por operador, dispara el rescale a través de Temporal y aísla cualquier falla a un solo job sin tocar al resto de la flota, todo mientras el equipo de plataforma sigue viendo un único punto de control para 30.000 pipelines distintos. Para sostener eso, Netflix modificó la recolección de métricas del JobManager de forma que soporte jobs de hasta 3.000 subtasks, sumó filtrado de métricas del lado del servidor, preservó los subgrafos conectados por FORWARD durante el rescalado y agregó manejo de backpressure en los sinks. Te puede servir nuestra cobertura sobre cómo elegir la herramienta correcta para automatizar despliegues, disponible en esta comparativa entre Jenkins y GitHub Actions.
El tema de las conexiones FORWARD no es menor. Cambiar el paralelismo a través de una conexión FORWARD puede forzar una redistribución completa de datos, algo que ya venía apareciendo en discusiones de desarrolladores de Apache Flink. La implementación de Netflix mantiene juntos los operadores conectados por FORWARD para evitar ese costo. Sigue habiendo un issue abierto en la comunidad sobre casos donde operadores con mucha carga pueden verse afectados por decisiones de escalado basadas en output ratio.
¿En qué se diferencia de otros autoscalers como KEDA?
KEDA y otros autoscalers “genéricos” event-driven escalan workloads a partir de métricas o eventos externos al proceso, como el tamaño de una cola o el uso de CPU del pod. El Flink Autoscaler, en cambio, razona sobre el grafo interno del dataflow y la capacidad real de cada operador.
Esa diferencia importa para stateful streaming. Un autoscaler externo no sabe que un operador con estado tiene un costo de rescalado distinto al de uno sin estado (recuperar terabytes de estado no es gratis). El autoscaler de Flink sí lo tiene en cuenta porque opera dentro del propio runtime del job.
¿Por qué Netflix usa un target de utilización de 0,45 en vez del 0,7 por defecto?
Netflix configura un utilization target de 0,45, por debajo del 0,7 que trae por default la comunidad de Flink, para reducir el rescalado agresivo de jobs stateful grandes. Un target más conservador deja margen antes de disparar un rescale, algo clave cuando cada rescale de un job con mucho estado implica mover y recuperar ese estado. En evitar caídas en infraestructuras que operan a esta escala profundizamos sobre esto.
¿Vale ese margen extra el costo de recursos ociosos? Para Netflix, sí, porque el costo de un rescale mal calculado en un job de 3.000 subtasks es mayor que el de correr un poco por debajo del límite óptimo de CPU.
¿Qué sigue para Netflix con Flink 2 y el estado desagregado?
Netflix planea migrar el resto de sus casos internos de autoscaling al proyecto open source y está investigando la arquitectura de estado desagregado de Flink 2 para atacar el costo de recuperación de estado durante el rescalado. Es la pieza que todavía le falta resolver del todo: aunque el nuevo autoscaler ya no sobreaprovisiona por operador, mover estado grande sigue siendo la parte cara de cualquier rescale.
Qué significa esto para equipos de datos en Latinoamérica
Pocas empresas de la región operan streaming a la escala de Netflix, pero el patrón de fondo aplica igual: un autoscaler cluster-level suele sobreaprovisionar para cubrir el operador más pesado del pipeline. Si tu equipo corre Flink o Kafka Streams sobre un clúster propio en AWS, GCP o infraestructura on-premise, vale la pena revisar si estás pagando de más por esa misma razón antes de sumar más nodos. Para la infraestructura web y de dominios de la empresa (no la parte de streaming), en Argentina donweb.com es una opción de hosting y VPS local.
Errores comunes al pensar el autoscaling de streaming
- Confundir el autoscaling de Kubernetes con el de Flink. Escalar pods vía HPA no resuelve el problema si el cuello de botella está en el paralelismo interno de un operador stateful; son dos capas distintas.
- Copiar el utilization target de 0,7 sin ajustarlo. Ese es el default comunitario, pero Netflix lo bajó a 0,45 para jobs grandes con mucho estado. Empezar conservador y ajustar con datos propios rinde mejor que copiar el número de otro.
- Ignorar las conexiones FORWARD al rescalar. Romper el paralelismo de subgrafos conectados por FORWARD dispara redistribuciones de datos caras que se podían evitar manteniendo esos operadores juntos.
- Activar el autoscaler sin instrumentar métricas por operador. Sin throughput y busy time reales por vértice, el cálculo de paralelismo queda a ciegas, con decisiones tan malas como las del viejo modelo cluster-level.
Preguntas Frecuentes
¿Qué es el Apache Flink Autoscaler?
Es un componente open source de Apache Flink que calcula el paralelismo necesario para cada operador de un job de streaming a partir de su throughput y busy time reales. Reemplaza los enfoques que escalan el clúster completo como una sola unidad.
¿Cómo ahorró Netflix 1,1 millones de dólares con Flink?
Un equipo interno de Netflix redujo su gasto anualizado en cómputo Flink un 58% al pasar del autoscaler cluster-level al nuevo autoscaler por operador, lo que equivale a cerca de USD 1,1 millones al año según InfoQ. El ahorro viene de dejar de sobreaprovisionar operadores livianos para cubrir a los pesados. Sobre eso hablamos en optimizar el alcance internacional de este tipo de proyectos.
¿Cuál es la diferencia entre un autoscaler a nivel de clúster y uno por operador?
Un autoscaler a nivel de clúster ajusta la cantidad total de TaskManagers y aplica la misma decisión a todos los operadores del job. Uno por operador, como el nuevo Flink Autoscaler, calcula el paralelismo de cada vértice del grafo por separado, algo más preciso en pipelines con branches y joins de distinto peso.
¿Qué es el proyecto DS2?
DS2 es la investigación de sistemas en la que se basa la técnica de estimación de paralelismo por operador del nuevo autoscaler. El equipo detrás del trabajo probó primero un análisis de ruta crítica más complejo antes de adoptar DS2 como baseline simple, que terminó funcionando mejor.
¿El Flink Autoscaler de Netflix funciona con el Flink Kubernetes Operator?
Netflix no despliega el autoscaler directamente vía el Flink Kubernetes Operator. Lo integró con un servicio propio en Spring Boot que usa workflows de Temporal para aislar las decisiones de autoscaling job por job dentro de su control plane interno.
Conclusión
El cambio de Netflix no es solo una optimización interna: es la confirmación pública de que el modelo cluster-level de autoscaling se quedó corto para streaming stateful a gran escala. El dato del 58% de ahorro en un equipo, sobre una flota de más de 30.000 jobs, sugiere que hay margen similar en cualquier organización que todavía escale Flink por clúster completo. Lo que sigue mirar de cerca es Flink 2 y su arquitectura de estado desagregado, la pieza que falta para bajar también el costo de recuperar estado en cada rescale.






