Netflix Migrates Over 30,000 Flink Jobs to an Open-Source Autoscaler, Citing a 58% Cost Cut for One Team
Netflix is consolidating two Flink autoscalers into the open-source community version across 30,000+ jobs, with one team reporting a 58% compute-cost cut worth $1.1 million a year.
Overview
Netflix is converging on the open-source Apache Flink Autoscaler across a fleet of more than 30,000 Flink jobs, phasing out an in-house autoscaler the company built around 2019, according to a Netflix Technology Blog post and InfoQ. One team at Netflix has cut its annualized Flink compute expenditure by 58%, saving approximately $1.1 million a year, according to both outlets.
What We Know
- Netflix has run stream processing on Apache Flink since 2017 and, as of 2026, operates more than 30,000 Flink jobs across multiple AWS regions, according to Netflix and InfoQ. Most of those jobs are generated by Netflix’s internal Data Mesh platform, while a smaller, growing set are custom jobs built by teams for use cases including personalization, advertising, and live events, according to Netflix.
- Netflix’s original autoscaler, built around 2019, ran on the company’s Mantis platform and consumed cluster-level telemetry — CPU, network utilization, Kafka lag, input rate, and consume rate — from its Atlas telemetry system, according to Netflix and InfoQ. Netflix says it “reliably cut resource usage by 25–45% across thousands of managed pipelines,” but the system could only resize an entire cluster by adjusting the total TaskManager count, moving every operator in a job together — a design that could not handle the multi-operator, stateful pipelines teams increasingly built for advertising, recommendations, and games, according to Netflix.
- The Apache Flink Autoscaler that Netflix is adopting instead reasons from inside each job rather than watching cluster-level containers from outside, calculating “required parallelism for individual vertices” by examining each operator’s processing rate, according to InfoQ. Netflix describes the underlying metric as each operator’s “true processing rate”: observed throughput divided by the fraction of time the operator is actually busy, extrapolated to what it could sustain at full utilization, according to Netflix.
- Rather than deploying the community autoscaler through the Flink Kubernetes Operator, Netflix built its own integration as a Spring Boot service orchestrated with Temporal workflows, running one long-running workflow per job, according to Netflix and InfoQ.
- To make the open-source autoscaler work at Netflix’s scale, engineers changed Flink’s JobManager to cache and clean up transient metric names instead of rescanning them on every fetch, and added server-side metric filtering. This let the autoscaler handle jobs with up to 3,000 subtasks, according to Netflix and InfoQ; Netflix says it had previously struggled above roughly 1,000 subtasks.
- Netflix runs the autoscaler at a target utilization of 0.45, below the Apache Flink community’s default of 0.7, “deliberately trading a little efficiency for stability,” according to Netflix.
- The open-source autoscaler reached general availability for Netflix’s custom jobs last year, according to Netflix. Netflix’s client telemetry and logging team, per that same post, is the one that achieved the reported 58% reduction in annualized Flink compute expenditure — a savings figure of approximately $1.1 million a year that InfoQ also reports.
What We Don’t Know
- Netflix has not disclosed what share of its 30,000-plus jobs have migrated to the open-source autoscaler versus the legacy Mantis-based system, or a target date for full consolidation.
- Netflix has not published cost or resource-usage figures for teams beyond the client telemetry and logging team cited in its blog post.
Analysis
Netflix frames the migration as a build-versus-buy case study rather than a purely technical one. The company says it built its first autoscaler in-house “when there was no mature option suited to our platform,” but is migrating to the community version because it can “scale workloads our homegrown system was never designed for,” according to Netflix. Looking ahead, Netflix says it plans to investigate Flink 2’s disaggregated state architecture — which keeps state in external storage rather than on local disk — to address what it now considers the largest remaining bottleneck in scaling stateful jobs: the restart and state-recovery process itself, according to Netflix and InfoQ.