Flink 作业的自动扩缩容,表面上都是“指标高了就扩容,指标低了就缩容”,真正落地时却会遇到两个完全不同的问题:一类方案理解 Flink 内部的反压、处理速率和并行度,另一类方案只负责根据外部指标调整 Kubernetes 工作负载。它们都能改变副本或并行度,但决策依据、调整粒度和风险边界并不相同。
本文把这两类方案称为“两个 Flink autoscaler”:Flink 原生的作业级自动扩缩容,以及 Kubernetes 层面的通用扩缩容器。重点不在于给出唯一答案,而在于帮助你判断,自动扩缩容究竟应该发生在 Flink 内部,还是发生在 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 指标;
minReplicaCount和maxReplicaCount根据下游容量设置;- 确认新增 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 调整目标工作负载规模”。生产环境还需要验证以下行为:
- TaskManager 新增后是否能被 JobManager 调度并使用;
- Flink 作业的算子并行度是否随资源变化而变化;
- 缩容时是否会触发不必要的状态迁移或任务重启;
- lag 降低后,下游系统是否仍处于安全状态;
- 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 更懂平台资源管理。选择时先定义要控制的对象,再选择能够看到相应信号的控制器,通常比从某个工具的默认配置开始更可靠。