如果你管理过大规模流处理集群,一定体会过这样的痛苦:业务方说某个任务延迟涨了,你手动把并行度调大,过两天又发现资源浪费严重,再调小。对于成百上千个任务,靠人肉调参根本忙不过来。Netflix 从 2017 年就开始跑 Apache Flink,2019 年自研了一套自动扩缩容(autoscaler)系统,运行在自研的 Mantis 平台上,靠 Atlas 里的集群级指标(CPU、网络利用率、Kafka lag 等)来调整整个集群的 TaskManager 数量,效果也不错,把数千条管道的资源消耗降低了 25% 到 45%。但问题是,这套旧系统把整个集群当成一个整体来缩放,所有算子(operator,也就是数据流处理逻辑里的一个处理节点,比如过滤、窗口聚合)共享同一个伸缩决策。

这在大规模有状态任务上就露馅了。Netflix 的生产任务里,很多是包含分支、join(合并多个流)和 TB 级状态的复杂数据流,不同阶段的计算负载差异巨大。比如一个算子可能 CPU 打满,而另一个算子却在空转,但集群级缩放只能一荣俱荣、一损俱损,没法针对瓶颈算子单独扩容。结果就是要么资源浪费,要么瓶颈没消除。

于是他们把目光转向了 Apache Flink 社区的开源 Autoscaler。Flink 自带的这个方案思路完全不同——它不再看集群宏观指标,而是利用任务运行时暴露的指标,估计每个算子的**真实处理速率(True Processing Rate)**。怎么估计?就是通过吞吐量和忙碌时间(busy time,指算子在处理数据而非空等的时间占比)算出一个“这个算子目前到底能处理多少条数据/秒”的真实能力。然后顺着任务图(job graph,数据流执行顺序的可视化结构)逐个计算每个顶点(vertex,即算子在物理执行时的单元)所需的并行度。这套设计记录在 Flink 的 FLIP-271(一个提案文档,Flink 社区用编号管理重大功能改动)里,专门解决异构流任务的自动缩放和状态应用的缩放成本问题。

这个技术路径还和波士顿大学一个叫 DS2 的研究项目有关。研究员 Vasiliki Kalavri 参与了这项工作,起初他们尝试用关键路径分析,后来发现直接用真实处理速率这个简单的基线就够了。“这个非常简单的想法居然效果极好”,后来就成了 Flink 自动缩放的核心。最让我意外的是,复杂问题有时不需要复杂解法,一个精心衡量的简单指标反而能穿透问题的本质。

Netflix 没有直接通过 Flink Kubernetes Operator 部署开源版,而是把它的核心逻辑整合进了自家内部的控制平面。具体做法:用一个 Spring Boot 服务,通过 Temporal 工作流(Temporal 是一个可靠的工作流编排引擎,能管理有重试、超时补偿的长期任务)来隔离每个任务的缩放决策,避免互相干扰。他们还改进了 JobManager 的指标采集,让它能支持最多 3000 个 subtask(子任务,并行度拆出来的每个并发实例)的大任务;加了服务端指标过滤,减少不必要的网络传输;缩容的时候保留 FORWARD 连接的子图(FORWARD 连接指上游算子直接把数据推给下游、不做重分区的连接方式);还专门处理了 sink(数据输出端)背压的情况。

这里有个值得注意的工程坑:FORWARD 连接在调整并行度时,如果两个算子原本是直接对接的,改变并行度可能就需要重新分发数据,成本很高。Netflix 的实现是让这种连在一起的算子保持连接、不拆开,这比 Flink 社区讨论里提到的做法更稳妥。另外还有一个悬而未决的问题,GitHub 上 FLINK-38538 指出,基于输出比的缩放决策可能影响忙碌算子,说明开源方案在极端场景下仍有待完善。

从效果看,Netflix 公布了一个团队的例子:把年度 Flink 计算支出降低了 58%,每年省下约 110 万美元。这不只是钱的问题,更重要的是运维效率——以前需要人工盯着的缩放决策,现在可以自动化且按算子粒度精确控制。和 KEDA 这类通用事件驱动缩放器(它依赖外部指标或事件来扩展工作负载)相比,Flink Autoscaler 的优势在于它理解内部数据流图和算子容量,而不是只看着外部排队长度瞎猜。

Netflix 目前把利用率目标设置在 0.45(即目标利用率为 45%),低于 Flink 社区默认的 0.7,为的是减少大状态任务被频繁缩放的扰动。他们计划把剩余的旧 autoscaler 用例全部迁移到开源版本,同时正在调研 Flink 2 的 disaggregated state(分离状态存储,把状态与计算分开管理)架构,以解决缩放时状态恢复成本高的问题。

对我们有什么启发?如果你的团队在用 Flink 做流处理,而且任务里混着有状态、多分支、大状态量的任务,别再自己去写集群级自动缩放——直接评估 Flink 自带的 Autoscaler,它在算子级粒度的处理上已经过 Netflix 这种规模的验证。即便不用 Flink,这个思路也通用:自动伸缩的单位要细化到真正有瓶颈的组件,而不是整个集群一刀切;衡量指标要选“真实处理能力”这种能反映实际吞吐的,而不是只看 CPU 或外部延迟。另外,注意背压、FORWARD 连接这类细节,缩容时优先保证数据通道不变,避免重建开销。迁移大状态任务时,先模拟调整并行度,观察状态恢复时间和资源消耗,必要时像 Netflix 那样把利用率目标调低来换稳定性。

Instrumentation at Scale: Having Your Performance Cake and Eating It Too
Instrumentation at Scale: Having Your Performance Cake and Eating It Too
阅读原文 → 返回 AI 技术文档

内容与图片版权归原作者所有 · 原文: https://www.infoq.com/news/2026/09/netflix-flink-autoscaler/?utm_campaign=infoq_content&utm_source=infoq&utm_medium=feed&utm_term=global