两个 Flink 自动扩缩容器的故事:如何选择适合你的弹性方案

2026-08-22 34 预计阅读时间: 1 分钟
来源: netflixtechblog.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:12 分钟

Flink 作业的自动扩缩容,表面上都是“指标高了就扩容,指标低了就缩容”,真正落地时却会遇到两个完全不同的问题:一类方案理解 Flink 内部的反压、处理速率和并行度,另一类方案只负责根据外部指标调整 Kubernetes 工作负载。它们都能改变副本或并行度,但决策依据、调整粒度和风险边界并不相同。

本文把这两类方案称为“两个 Flink autoscaler”:Flink 原生的作业级自动扩缩容,以及 Kubernetes 层面的通用扩缩容器。重点不在于给出唯一答案,而在于帮助你判断,自动扩缩容究竟应该发生在 Flink 内部,还是发生在 Flink 之外。

两种方案解决的是不同问题

Flink 原生 autoscaler 通常围绕作业拓扑和算子运行状态做决策。它可以观察处理速率、输入速率、忙碌时间、反压以及算子级别的负载,并据此推算某个算子需要多少并行度。

这种方式适合以下场景:

  • 作业是持续运行的流处理任务;
  • 瓶颈集中在某些算子,而不是整个集群;
  • 希望调整的是 Flink 作业并行度;
  • 需要结合 checkpoint、状态和拓扑信息做变更;
  • 能够接受扩缩容过程中的重启、状态恢复或重新部署成本。

它的核心优势是“知道自己在扩什么”。例如,窗口聚合、连接或下游 sink 可能分别处于不同负载水平,仅凭一个 CPU 指标很难判断应该增加哪个算子的并行度。作业级 autoscaler 可以把决策建立在 Flink 的处理语义和运行指标之上。

Kubernetes 通用扩缩容:管理工作负载的数量

另一类方案运行在 Kubernetes 层,典型实现包括 HPA、KEDA 或自定义 controller。它们通过 CPU、内存、Prometheus 指标、Kafka lag 等外部信号,调整 Deployment、StatefulSet 或其他工作负载的规模。

这类方案适合的对象通常是:

  • 能够独立运行的 Flink JobManager 或 TaskManager 工作负载;
  • 一批短任务或按需提交的 Flink 作业;
  • 主要瓶颈可以用队列长度、消费延迟或资源使用率表达;
  • 团队已经有成熟的 Kubernetes 监控与发布体系。

它的优点是接入简单、平台统一。缺点也很明确:Kubernetes 知道 Pod 的资源使用情况,却未必知道 Flink 拓扑中哪个算子正在积压,也不知道改变 TaskManager 数量是否真的会改变作业的有效并行度。

换句话说,HPA 或 KEDA 可以回答“需要多少 Kubernetes 资源”,但不一定能回答“Flink 作业应该把并行度加到哪里”。

选择时要看四个维度

1. 扩缩容粒度

如果目标是调整算子并行度,Flink 原生 autoscaler 更贴合问题。它可以针对作业拓扑中的实际瓶颈做决策。

如果目标是扩展一组相对独立的工作负载,或者只是根据外部队列长度增加处理实例,Kubernetes 层扩缩容通常更直接。

需要特别注意:增加 TaskManager Pod 不一定自动增加某个算子的并行度。作业是否会使用新增资源,取决于部署模式、资源配置和 Flink 作业本身的调度方式。

2. 扩缩容触发信号

CPU 使用率是一个容易获得的指标,但不一定是流处理系统的好指标。一个算子可能 CPU 使用率不高,却因为外部系统响应慢而出现反压;也可能 CPU 使用率很高,但输入流量马上就会下降。

更适合 Flink 的信号通常包括:

  • 输入速率与实际处理速率;
  • 算子 busy time;
  • 反压状态;
  • Kafka consumer lag 或其他队列积压;
  • checkpoint 时长与失败率;
  • sink 端延迟和错误率。

实际使用时,最好把业务延迟目标和资源指标结合起来。例如,Kafka lag 适合作为外部扩缩容信号,但仍需要确认增加并行度不会把压力转移到数据库、API 或消息系统。

3. 状态与变更成本

有状态 Flink 作业的扩缩容并非简单地增加几个进程。作业可能需要触发 savepoint、重新部署、重新分配状态,并等待 checkpoint 或恢复完成。状态规模越大,扩缩容动作的代价越高。

因此,自动扩缩容必须包含:

  • 扩缩容冷却时间;
  • 最小和最大并行度;
  • checkpoint 或 savepoint 策略;
  • 失败回滚;
  • 变更后的观测窗口;
  • 防止扩缩容来回振荡的稳定机制。

如果一次扩缩容需要数分钟才能完成,控制器就不应该每几十秒做一次新的决策。

4. 依赖系统的容量边界

Flink 只是数据链路的一环。扩容前需要检查下游系统是否承受得住新增吞吐:数据库连接池、HTTP API 限流、Kafka 分区数、对象存储请求速率,都可能成为新的瓶颈。

自动扩缩容的上限不应只由 Flink 资源预算决定,也应由最脆弱的依赖系统决定。一个实用做法是为每个作业记录“可安全处理的最大吞吐”,并把它转换为并行度上限或外部扩缩容上限。

一个可改造的 Kubernetes 示例

下面是一个简化的 KEDA ScaledObject 示例。它假设 Flink 作业从 Kafka 读取数据,团队已经通过 Prometheus 暴露了一个代表消费积压的指标。示例展示的是 Kubernetes 层的自动扩缩容思路,不代表所有 Flink 部署模式都可以直接使用。

运行前需要修改以下内容:

  • scaleTargetRef.name 改成实际的工作负载名称;
  • Prometheus 地址改成集群内可访问的地址;
  • 指标查询改成实际的 lag 指标;
  • minReplicaCountmaxReplicaCount 根据下游容量设置;
  • 确认新增 Pod 会被 Flink 作业实际使用。
apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
  name: flink-taskmanager-scaler
  namespace: streaming
spec:
  scaleTargetRef:
    name: flink-taskmanager
  pollingInterval: 30
  cooldownPeriod: 300
  minReplicaCount: 2
  maxReplicaCount: 12
  advanced:
    horizontalPodAutoscalerConfig:
      behavior:
        scaleUp:
          stabilizationWindowSeconds: 60
        scaleDown:
          stabilizationWindowSeconds: 600
  triggers:
    - type: prometheus
      metadata:
        serverAddress: http://prometheus.monitoring.svc:9090
        metricName: kafka_consumer_lag
        threshold: "10000"
        query: |
          sum(kafka_consumergroup_lag{namespace="streaming",group="orders-flink"})

这个配置只解决了“根据 lag 调整目标工作负载规模”。生产环境还需要验证以下行为:

  1. TaskManager 新增后是否能被 JobManager 调度并使用;
  2. Flink 作业的算子并行度是否随资源变化而变化;
  3. 缩容时是否会触发不必要的状态迁移或任务重启;
  4. lag 降低后,下游系统是否仍处于安全状态;
  5. Prometheus 指标暂时缺失时,控制器是否会采取保守动作。

如果这些问题无法得到明确答案,直接把 Kubernetes 扩缩容器接到 Flink 工作负载上,可能只是在改变 Pod 数量,而没有真正改善作业吞吐。

更稳妥的组合方式

在复杂平台中,两类 autoscaler 可以分工,但不应同时无边界地控制同一个变量。

一种可行的职责划分是:

  • Flink 原生 autoscaler 负责作业并行度;
  • Kubernetes autoscaler 负责集群或节点资源;
  • Cluster Autoscaler 负责节点数量;
  • 平台控制器负责发布、savepoint 和失败回滚。

这相当于把控制环拆开:Flink 根据作业负载决定需要多少并行处理能力,Kubernetes 根据资源请求决定需要多少 Pod,云平台再根据 Pod 调度结果决定需要多少节点。

拆分控制环时要避免多个控制器同时修改同一项配置。例如,Flink autoscaler 正在把算子并行度从 4 调整到 8,另一个 HPA 又根据 CPU 把 TaskManager 数量从 3 调整到 10,两个动作叠加后可能造成资源浪费、状态迁移或调度震荡。

部署前应明确每个控制器的:

  • 输入指标;
  • 控制目标;
  • 最小与最大边界;
  • 冷却时间;
  • 失败处理方式;
  • 变更审计记录。

落地清单

开始启用自动扩缩容前,可以用下面的顺序做验证:

  • 先记录基线:输入速率、处理速率、端到端延迟、checkpoint 时长和资源使用率;
  • 找出真正的瓶颈算子,而不是只看整体 CPU;
  • 为 Kafka、数据库、API 和 sink 设置容量上限;
  • 通过压测确认扩容后吞吐确实提高;
  • 设置扩缩容冷却时间和并行度边界;
  • 在测试环境验证 savepoint、状态恢复和缩容行为;
  • 观察至少一个完整业务高峰,再决定是否放宽最大并行度;
  • 为指标缺失、部署失败和恢复超时准备保守策略。

两个 Flink autoscaler 的差异,最终不在于哪个组件名称更流行,而在于它们观察系统的层次不同。作业级 autoscaler 更懂 Flink 的处理逻辑,Kubernetes 层 autoscaler 更懂平台资源管理。选择时先定义要控制的对象,再选择能够看到相应信号的控制器,通常比从某个工具的默认配置开始更可靠。


相关推荐