A Tale of Two Flink Autoscaler... Note

A Tale of Two Flink Autoscalers

Netflix currently runs two Flink autoscalers, aiming to consolidate to one open-source solution. They initially built an in-house autoscaler to manage thousands of Flink jobs, which proved effective but limited in handling complex, multi-operator jobs. This homegrown system relied on external metrics, which could be inaccurate or miss internal job issues. The Apache Flink community later developed a more sophisticated autoscaler that reasons from within the job itself. This second autoscaler estimates an operator's true processing rate to determine optimal parallelism per vertex. Netflix is converging on this open-source solution due to its ability to scale stateful jobs and allow granular configuration. Implementing the OSS autoscaler at Netflix required significant adaptation, including integrating it into their existing control plane. They also made modifications to Flink's runtime for better metric collection and to handle specific job structures like forward-connected subgraphs and sink limitations. By adopting the OSS autoscaler, Netflix has achieved substantial cost savings and improved resource utilization. Lessons learned emphasize that metric quality is crucial, tunable defaults are beneficial, and adopting existing solutions before extending them is a wise strategy. They plan to fully migrate to the OSS-based autoscaler to streamline operations.
CdXz5zHNQW_HJCpS7bDND.png