Lyft 将数百个生产 Flink 作业从一套始建于 2020 年的内部 Kubernetes Operator,迁移到了 Apache Flink Kubernetes Operator。改变的不只是控制器代码:统一到社区控制面后,团队能够在整个作业集群中推进最近状态升级、原地自动扩缩容和资源自动调优。
对于维护大规模流处理平台的团队,这类迁移的核心价值在于减少自研控制面的长期负担,同时把升级、状态恢复和资源治理变成声明式能力。
真正需要迁移的是控制面语义
从一个 Operator 切换到另一个 Operator,并不等同于替换一组 Kubernetes YAML。自研控制器通常已经隐含了大量平台约定:
- 作业如何提交、挂起和恢复;
- Checkpoint、Savepoint 与高可用元数据保存在哪里;
- JobManager 失败后由谁重建;
- Flink 版本或业务 JAR 更新时采用哪种升级策略;
- CPU、内存和并行度由应用团队还是平台团队决定;
- 告警、审计、回滚和发布状态如何接入现有系统。
因此,迁移时应建立“旧控制器行为到新 CRD 字段”的映射,而不是直接翻译资源定义。尤其要确认以下三类状态不会丢失:
- 业务状态:窗口、聚合、连接器 offset 等 Flink 状态。
- 控制面状态:作业版本、部署阶段、升级中的中间状态。
- 外部系统状态:Kafka 消费进度、输出端事务以及对象存储中的 Checkpoint。
这也是最近状态升级值得关注的原因。与每次升级都显式触发 Savepoint 相比,last-state 模式可以利用最近一次可恢复状态完成升级,缩短部分作业的升级路径。但它依赖可靠的高可用元数据和持久化 Checkpoint,不能被理解为“无需状态存储”。
用 FlinkDeployment 描述一个有状态作业
下面是一个可以直接改造的最小示例。运行前需要完成三项修改:把镜像替换为包含业务 JAR 的镜像,把 jarURI 改为镜像内的实际路径,并将对象存储地址替换为可访问的位置。集群中还必须已经安装与目标 Flink 版本兼容的 Apache Flink Kubernetes Operator。
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
name: orders-stream
namespace: streaming
spec:
image: registry.example.com/streaming/orders-job:1.0.0
flinkVersion: v1_18
flinkConfiguration:
taskmanager.numberOfTaskSlots: "2"
high-availability.type: kubernetes
high-availability.storageDir: s3://my-flink-state/ha/orders-stream
state.checkpoints.dir: s3://my-flink-state/checkpoints/orders-stream
execution.checkpointing.interval: "60s"
kubernetes.operator.job.autoscaler.enabled: "true"
job.autoscaler.scaling.enabled: "true"
job.autoscaler.stabilization.interval: "2m"
job.autoscaler.metrics.window: "10m"
job.autoscaler.target.utilization: "0.70"
job.autoscaler.target.utilization.boundary: "0.10"
serviceAccount: flink
jobManager:
resource:
cpu: 1
memory: 2048m
taskManager:
resource:
cpu: 1
memory: 2048m
job:
jarURI: local:///opt/flink/usrlib/orders-job.jar
parallelism: 4
upgradeMode: last-state
state: running
保存为 orders-stream.yaml 后,可以这样部署和观察:
kubectl create namespace streaming --dry-run=client -o yaml | kubectl apply -f -
kubectl apply -f orders-stream.yaml
kubectl -n streaming get flinkdeployment orders-stream -w
检查 Operator 是否接受了配置:
kubectl -n streaming describe flinkdeployment orders-stream
kubectl -n streaming get pods
kubectl -n streaming logs deployment/flink-kubernetes-operator --tail=200
这个示例展示的是配置入口,不代表任何版本组合都能直接完成原地扩缩容。是否能够避免完整重启,取决于 Flink、Operator、调度模式和自动扩缩容组件的兼容关系。上线前应针对实际版本验证支持矩阵。
自动扩缩容不是简单修改副本数
无状态服务通常可以通过增加 Pod 副本扩大吞吐量,但 Flink 作业的并行度与状态分区绑定。一次扩缩容可能触发状态重新分布、网络数据交换和短期反压,因此控制器需要考虑的不只是当前 CPU 使用率。
更实用的策略通常会同时观察:
- Source backlog 或消费延迟;
- 算子繁忙度与空闲时间;
- 反压持续时间;
- Checkpoint 时长与失败率;
- 扩缩容后的状态搬迁成本;
- 目标利用率以及稳定窗口。
示例中的稳定时间和指标窗口用于降低频繁抖动。生产环境还应给并行度设置上下限,并对大状态作业采用更保守的策略。否则,扩容刚完成就缩容,反复发生的状态重分布可能比少量资源闲置更昂贵。
资源自动调优也不应只理解为“自动减小内存”。比较稳妥的做法是先在建议模式中收集数据,让系统给出 CPU、内存和并行度建议,再对低风险作业启用自动执行。状态后端、序列化格式、连接器缓冲区和堆外内存都会影响调优结果,单一资源指标不足以覆盖这些因素。
数百个作业应如何分批迁移
大规模迁移更适合按风险分层,而不是按团队一次性切换。
1. 建立作业清单
至少记录 Flink 版本、状态大小、Checkpoint 时长、外部依赖、当前升级模式、恢复时间目标和业务重要性。无状态或状态较小的作业适合作为第一批候选。
2. 固化迁移前状态
对关键作业,在切换控制器前创建经过验证的 Savepoint,并确认新环境能够读取对应的存储路径、凭据和状态后端格式。last-state 适合日常升级,但跨控制面迁移时,显式恢复点通常更容易审计和回滚。
3. 避免两个 Operator 同时管理同一作业
旧控制器和新控制器不能同时对同一工作负载执行协调。迁移流程应包含明确的所有权交接,例如:
# 示例流程,资源名称需按实际环境修改
kubectl -n streaming scale deployment/internal-flink-operator --replicas=0
kubectl -n streaming apply -f orders-stream.yaml
kubectl -n streaming get flinkdeployment orders-stream -w
如果旧 Operator 还管理其他作业,就不能直接整体缩容。可以通过 namespace、标签选择器或分片部署隔离管理范围,再逐批转移所有权。
4. 用真实故障验证恢复能力
除了检查作业进入 RUNNING,还应在预发布环境测试 JobManager 重建、TaskManager 丢失、对象存储短暂不可用、升级失败和回滚。只有恢复后的输出语义、消费 offset 和延迟均符合预期,迁移才算完成。
采用前的检查清单
Apache Flink Kubernetes Operator 可以减少自研控制器的维护成本,但不会自动解决所有流处理问题。落地前建议确认:
- Operator、Flink 和 Kubernetes 版本已经过兼容性验证;
- Checkpoint、Savepoint 和 HA 元数据位于持久化存储;
last-state、savepoint与无状态升级分别有明确使用条件;- 自动扩缩容设置了并行度边界、稳定窗口和告警;
- 迁移期间只有一个控制器拥有作业;
- 每个迁移批次都有停止条件和回滚路径;
- 平台监控覆盖 Operator 协调失败,而不只是 Flink 作业指标。
Lyft 的迁移说明,成熟的社区 Operator 不只是部署工具,更可以成为整个流处理平台的统一控制面。真正的收益来自标准化之后的规模效应:同一套升级策略、扩缩容规则和资源治理机制能够覆盖数百个作业,而平台团队不再需要独自维护一套不断追赶 Flink 演进速度的内部控制器。