欢迎光临
我们一直在努力

FlinkK8sOperator源码深潜(下):JobManager与TaskManager Controller如何生成K8s集群资源

FlinkK8sOperator源码深潜(下):JobManager与TaskManager Controller如何生成K8s集群资源

【免费下载链接】flinkk8soperator Kubernetes operator that provides control plane for managing Apache Flink applications 【免费下载链接】flinkk8soperator 项目地址: https://gitcode.com/gh_mirrors/fl/flinkk8soperator

FlinkK8sOperator 是一个 Kubernetes Operator,作为 Apache Flink 应用的控制平面,让你用一个 FlinkApplication CR 就能自动拉起整个 Flink 集群。在上篇讲完状态机之后,本篇深入 pkg/controller/flink/ 目录,拆解 JobManager 与 TaskManager Controller 是如何把一份 CRD 声明"翻译"成 Deployment、Service、Ingress 等真实 K8s 资源的 👀。

一张图看懂:CR 到 K8s 资源的翻译链路

整个资源生成由状态机驱动,核心入口是 Controller.CreateCluster,它按顺序调用两个子控制器:

FlinkApplication CR


Controller.CreateCluster() ← flink.go#L260
├── jobManager.CreateIfNotExist() → Deployment + 2×Service (+Ingress)
└── taskManager.CreateIfNotExist() → Deployment

对应代码在 flink.go:先创建 JobManager 侧资源,再创建 TaskManager 侧资源,任一步失败都会通过 EventRecorder 往 CR 上打 Warning 事件,方便你在 kubectl describe 里直接看到失败原因。

FlinkK8sOperator 默认与InPlace双状态机,展示JobManager与TaskManager资源创建的状态流转

JobManager Controller:不止是 Deployment

一次调用创建 4 类资源

JobManagerController.CreateIfNotExist 位于 job_manager_controller.go,它一次会创建:

资源命名规则作用
Deployment {app}-{hash}-jm 拉起 JobManager Pod,副本数默认 1
通用 Service {app} 稳定入口,供 Ingress 暴露 Flink UI
版本化 Service {app}-{hash} 每版本一个,TaskManager 通过它精确连上对应版本的 JobManager
Ingress {app} 可选,配置 flinkIngressURLFormat 后生效

这里最精巧的是"版本化 Service"的设计:蓝绿部署时新旧两版集群并存,TaskManager 需要知道"我该连哪个 JobManager"。源码通过 InjectOperatorCustomizedConfig 把 jobmanager.rpc.address 注入到容器的 FLINK_PROPERTIES 环境变量中,指向 VersionedJobManagerServiceName 生成的版本化 Service 名——这就是 Flink 集群寻址的"最后一公里"。

容器的关键细节

FetchJobManagerContainerObj(job_manager_controller.go#L279-L318)构建了 JobManager 容器:

  • 就绪探针:每 5 秒 HTTP GET 一次 UI 端口的 /overview,JobManager Web UI 能响应才算就绪,Service 只会在真正可用后才开始转发流量;
  • 5 个标准端口:rpc(6123)、blob(6125)、query(6124)、ui(8081)、metrics(50101),均可在 CR 中覆盖;
  • 默认资源:4 CPU / 3072Mi,未指定 resources 时自动兜底;
  • Recreate 策略:JobManager 不能双副本并存,所以 Deployment 使用 RecreateDeploymentStrategyType,避免新旧副本抢端口。

TaskManager Controller:副本数不是拍脑袋定的

TaskManager 侧逻辑更简单——只创建一个 Deployment(task_manager_controller.go#L74-L92),但副本数的计算是 Flink 调度语义的体现:

replicas = ceil(parallelism / taskmanager.numberOfTaskSlots)

见 ComputeTaskManagerReplicas:你要跑并行度 24 的作业、每个 TaskManager 开 16 个 slot(默认值,见 config.go),则自动拉起 2 个副本。你只声明业务语义,资源数量由 Operator 换算。

TaskManager 容器(FetchTaskManagerContainerObj)与 JobManager 相比有两个明显差异:

  • 没有就绪探针:TaskManager 不需要对外提供服务,心跳由 Flink 自身管理;
  • 多一个 taskmanager.host 配置:注入为 $HOST_IP(Pod IP),让 JobManager 能反连 TaskManager 拉取指标。

两个"哈希":版本管理与幂等基石

源码里有两个容易混淆的 8 位字符串,都定义在 container_utils.go:

1. 应用哈希 HashForApplication:对 JobManager + TaskManager 两个 Deployment 模板做确定性 JSON 哈希,再取 FNV-32。任何 spec 变更都会产生新哈希 → 新的资源名 → 触发新版本集群,这是"改 CR 自动发版"的基石。

2. 随机 Pod 选择器 RandomPodDeploymentSelector:每次生成 Deployment 时随机产生,同时写入 Pod 标签和 Deployment Selector。注释里写明了动机:原地扩缩容时哈希可能不变,用随机串才能可靠区分"当前这一次创建"的 Pod。

另外所有资源都设置了 OwnerReferences 指向 FlinkApplication CR——删除 CR 时 K8s 垃圾回收机制会自动清理整条资源链,不留孤儿。

内存与配置如何进入 flink-conf.yaml

一个常见疑问:CR 里写的 taskmanager.memory.process.size 到底怎么生效?答案是不走 ConfigMap,走环境变量。renderFlinkConfig 把所有 Flink 配置(含按 Flink 版本分叉的内存参数:1.11 以下用 *.heap.size,1.11+ 用 *.memory.process.size)渲染成 YAML 文本,塞进容器的 FLINK_PROPERTIES 环境变量,Flink 启动脚本会自动解析它覆盖内置 flink-conf.yaml。

注意一个版本守卫:对 1.11 以下版本使用 systemMemoryFraction、对 1.11+ 使用 offHeapMemoryFraction 都会直接报错,从源头上避免用户写错配置项。

蓝绿部署:给资源名加上"版本号"

切换到 blueGreen 部署模式后,命名规则会多出一段:

{app}-{hash}-{version}-jm / {app}-{hash}-{version}-tm

见 getJobManagerName 与 getTaskManagerName。新旧两个版本的 Deployment、Service 可以安全共存:新版集群建好后双跑验证(DualRunning),确认通过再拆旧版。其状态流转如下:

FlinkK8sOperator 蓝绿部署状态机,展示新版本集群从创建到双跑验证的完整流程

小结:一套值得借鉴的 Controller 设计

回到源码目录,各文件的职责非常清晰:

  • 资源生成:job_manager_controller.go、task_manager_controller.go
  • 命名/标签/哈希/配置注入:container_utils.go
  • Flink 配置渲染与默认值:config.go
  • UI 入口:ingress.go
  • 集群装配与 Flink REST 调用:flink.go
  • CRD 类型定义:types.go

FlinkK8sOperator 给新手的最大启示是:Controller 的价值不在"调 API 建资源",而在于把业务语义(并行度、版本、蓝绿)编译成 K8s 原语,并用哈希 + OwnerReference + 幂等创建(CreateIfNotExist 对 AlreadyExists 的容忍处理)保证反复 Reconcile 时的稳定性。理解了这条链路,再去看任何 K8s Operator 源码都会顺畅很多。

【免费下载链接】flinkk8soperator Kubernetes operator that provides control plane for managing Apache Flink applications 【免费下载链接】flinkk8soperator 项目地址: https://gitcode.com/gh_mirrors/fl/flinkk8soperator

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

赞(0)
未经允许不得转载:171主机测评 » FlinkK8sOperator源码深潜(下):JobManager与TaskManager Controller如何生成K8s集群资源
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址