FlinkK8sOperator快速入门:5步部署你的第一个Flink流式作业(附完整配置解析)
2026/8/28 9:58:09 网站建设 项目流程

FlinkK8sOperator快速入门:5步部署你的第一个Flink流式作业(附完整配置解析)

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

FlinkK8sOperator 是一个 Kubernetes Operator,为运行在 Kubernetes 上的 Apache Flink 流式应用提供完整的控制平面:创建集群、提交作业、自动升级、保存点(Savepoint)与状态恢复统统由它托管。本文带你用 5 步在自有集群上部署第一个 Flink 流式作业,并逐字段解析 FlinkApplication 自定义资源的关键配置。

一、FlinkK8sOperator 是什么?

传统方式在 Kubernetes 上跑 Flink,需要你手动管理 JobManager、TaskManager、Service、Ingress 等一系列资源,升级和扩缩容更是繁琐。FlinkK8sOperator 通过 Kubernetes自定义资源(CRD)把这一切抽象成一个声明式的FlinkApplication对象:

  • 你只需提交一份 YAML 描述"我要什么"
  • Operator 持续监听该资源,自动创建 JobManager Deployment、TaskManager Deployment、Service,以及可选的 Ingress(Flink Web UI)
  • 应用更新时,Operator 自动完成"保存点 → 滚动切换 → 状态恢复",支持DualBlueGreen两种部署模式

上图展示了 Dual 部署模式下的状态机:从New/UpdatingClusterStartingSavepointingSubmittingJob,最终到达Running;失败时会进入RollingBackDeployFailed分支。

环境要求

组件版本要求
Kubernetes>= 1.10(< 1.13 需开启--feature-gates=CustomResourceSubresources=true
Apache Flink>= 1.7
kubectl已配置好集群凭证

二、准备工作:构建 Flink 应用镜像

Operator 启动 Flink 集群依赖一个 Docker 镜像,镜像中需要同时包含 Flink 运行时和你的应用代码(jar 包)。项目内置了可直接参考的 WordCount 示例:

  • 示例代码与 Dockerfile:examples/wordcount/
  • 示例自定义资源:examples/wordcount/flink-operator-custom-resource.yaml
  • 应用镜像编写规范:examples/README.md

💡 要点:镜像中的 jar 必须位于web.upload.dir指定的目录(示例中为/opt/flink),因为 Operator 是通过 JobManager 的 REST API 提交作业的。

三、5 步部署你的第一个 Flink 流式作业

第 1 步:创建 CRD、命名空间与 RBAC 权限

Operator 依赖自定义资源定义和一组最小权限角色,项目deploy/目录提供了全套清单:

kubectl apply -f deploy/crd.yaml kubectl apply -f deploy/namespace.yaml kubectl apply -f deploy/role.yaml kubectl apply -f deploy/role-binding.yaml

第 2 步:创建 Operator 配置 ConfigMap

Operator 的运行参数通过 ConfigMap 注入,模板见deploy/config.yaml

data: config: |- operator: ingressUrlFormat: "{{$jobCluster}}.{ingress_suffix}" logger: level: 4

{ingress_suffix}替换为你集群的 Ingress 域名后缀,即可让每个 Flink 应用拥有独立的 Web UI 地址。如果不需要 Ingress,留空即可,Operator 将不会创建 Ingress 资源。

第 3 步:部署 Operator

kubectl apply -f deploy/flinkk8soperator.yaml

确认 Operator Pod 处于RUNNING状态:

kubectl get pods -n flink-operator kubectl logs {pod-name} -n flink-operator

第 4 步:提交第一个 FlinkApplication

以内置 WordCount 示例为例(将image换成你自己的镜像地址):

kubectl apply -f examples/wordcount/flink-operator-custom-resource.yaml

示例核心内容(完整文件见examples/wordcount/flink-operator-custom-resource.yaml):

spec: image: flink-wordcount flinkVersion: "1.16" jarName: "wordcount-operator-example-1.0.0-SNAPSHOT.jar" entryClass: "org.apache.flink.WordCount" parallelism: 3 taskManagerConfig: taskSlots: 3 jobManagerConfig: replicas: 1

第 5 步:验证作业状态

Operator 观察到新资源后会自动拉起 Flink 集群并提交作业:

# 查看自动创建的 Deployment kubectl get deployments -n flink-operator # 查看自定义资源状态,phase 变为 Running 即成功 kubectl get flinkapplication.flink.k8s.io -n flink-operator wordcount-operator-example -o yaml # 查看 Operator 事件(创建集群、提交作业等) kubectl describe flinkapplication.flink.k8s.io -n flink-operator wordcount-operator-example

status中会返回集群健康度(health: Green)、可用 slot 数、作业 ID、checkpoint 计数等信息——这些信息都由状态机自动维护,你无需轮询任何 API。

四、FlinkApplication 完整配置解析

完整字段参考文档见docs/crd.md,类型定义位于pkg/apis/app/v1beta1/types.go。以下是最常用的字段速查表:

必填字段

字段说明
spec.image应用镜像,格式registry/repository[:tag]
spec.jarName镜像中待运行的 jar 包文件名
spec.parallelism作业级并行度
spec.flinkVersionFlink 版本,必须与镜像内版本一致
spec.jobManagerConfig.replicasJobManager 副本数(多副本需自行配置 HA)
spec.taskManagerConfig.taskSlots每个 TaskManager 的 slot 数

常用可选字段

字段说明
spec.flinkConfig透传给 Flink 的配置,如 checkpoint/savepoint 目录
spec.programArgs传给作业的启动参数(输入输出源等)
spec.deploymentMode更新部署模式:Dual(默认)或BlueGreen(零停机)
spec.scaleMode并行度变更策略:NewCluster(默认)或InPlace(实验性)
spec.deleteMode删除资源时的清理方式:Savepoint(默认)/ForceCancel/None
spec.restartNonce修改该值即可强制重启集群
spec.savepointPath从指定保存点恢复应用状态
spec.volumes/volumeMounts为 Job/TaskManager Pod 挂载卷

示例中的 checkpoint 配置建议保留,它是自动升级与故障恢复的基石:

flinkConfig: state.checkpoints.dir: file:///checkpoints/flink/externalized-checkpoints state.savepoints.dir: file:///checkpoints/flink/savepoints

两种部署模式怎么选?

  • Dual(默认):适合实时处理应用。Operator 先拉起新集群,再将旧作业带保存点取消、删除旧集群、在新集群恢复作业——停机时间极短。
  • BlueGreen:适合要求零停机的应用。新版本(蓝/绿)与旧版本并行运行,进入DualRunning阶段后,你通过设置spec.tearDownVersionHash手动确认下线旧版本。

两种模式的状态机细节见docs/state_machine.md

五、日常运维速查

操作命令
更新作业(自动带保存点滚动)修改 YAML 后kubectl apply -f <file>
查看状态与事件kubectl describe flinkapplication.flink.k8s.io <name>
强制重启集群修改spec.restartNonce后 apply
删除应用(默认先打保存点再下线)kubectl delete flinkapplication.flink.k8s.io <name>
本地开发调试 Operator参考docs/local_dev.md
更新后回滚设置spec.forceRollback: true(仅非 Running 阶段生效)

六、总结与下一步

通过本文的 5 步,你已经完成了 FlinkK8sOperator 的安装并部署了第一个 Flink 流式作业:CRD 与 RBAC → Operator 配置 → Operator 部署 → 提交 FlinkApplication → 验证 Running 状态。

接下来建议按此路径深入:

  • 自定义资源完整字段:docs/crd.md
  • 状态机各状态含义:docs/state_machine.md
  • 进阶用法(挂载卷、删除模式、回滚):docs/user_guide.md
  • 更多应用示例:examples/(含 Beam Python 示例)

把 Flink 应用的托管交给 Operator,你可以专注于流式逻辑本身,而不是 Kubernetes 的运维细节。

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

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询