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 自动完成"保存点 → 滚动切换 → 状态恢复",支持Dual和BlueGreen两种部署模式
上图展示了 Dual 部署模式下的状态机:从
New/Updating到ClusterStarting、Savepointing、SubmittingJob,最终到达Running;失败时会进入RollingBack与DeployFailed分支。
环境要求
| 组件 | 版本要求 |
|---|---|
| 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-examplestatus中会返回集群健康度(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.flinkVersion | Flink 版本,必须与镜像内版本一致 |
spec.jobManagerConfig.replicas | JobManager 副本数(多副本需自行配置 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),仅供参考