Netflix正在将其30000多个流式作业迁移至开源的Apache Flink Autoscaler,覆盖多个AWS区域。此前Netflix发现,其基于集群级别的自动扩缩容方案对包含不同处理需求算子的复杂有状态管道效果欠佳。Netflix表示,某团队通过该方案将Flink计算支出年化降低58%,每年节省约110万美元。

Netflix自2017年开始运行Apache Flink,并于2019年左右构建首个自动扩缩容系统。该系统运行在Mantis上,读取来自Atlas的集群级遥测数据,包括CPU、网络利用率、Kafka延迟、输入速率和消费速率,通过调整TaskManager总数量,在数千个管道上将资源使用量降低了25%至45%。其局限在于扩缩容粒度:由于基于集群而非单个算子决策,作业中所有算子共享同一扩缩容决策,对包含分支、连接操作和数TB状态的有状态管道越来越不适用。

Apache Flink Autoscaler利用运行作业暴露的指标,结合吞吐量和繁忙时间估算每个算子的真实处理速率,再遍历作业图,为各顶点分别计算所需并行度,该方法在FLIP-271中有详细描述。这项技术还借鉴了DS2项目的研究成果,参与该工作的系统研究员Vasiliki Kalavri表示,项目最初研究关键路径分析,随后采用真实处理速率作为更简单的基线方案,效果很好。

Netflix将该自动扩缩容器与内部控制平面集成,而非直接通过Flink Kubernetes Operator部署。一个Spring Boot服务使用Temporal工作流隔离各作业的扩缩容决策。Netflix还修改了JobManager的指标收集机制,可支持多达3000个子任务的作业,并增加服务端指标过滤、扩缩容时保留前向连接子图以及对Sink回压的处理逻辑。

与KEDA等基于外部指标或事件扩缩工作负载的通用事件驱动自动扩缩容器不同,Flink的自动扩缩容器基于内部数据流图和算子容量进行推理。Netflix目前使用的利用率目标为0.45,低于Flink社区默认的0.7,目标是减少大型有状态作业的激进式重新扩缩容。Netflix计划将剩余内部自动扩缩容实例迁移至开源实现,同时研究Flink 2的分离式状态架构,以解决重新扩缩容期间状态恢复的成本问题。