Skip to content
Charles Shao
Go back

Flink 深挖 · 部署模式与 Flink on K8s

–views

前面六站——架构、时间语义与窗口、端到端 Exactly-Once、DataStream、状态与容错与反压调优——都在本地的 docker / Session 集群里把作业调到了「跑得对、跑得快」,可那毕竟是一台笔记本上的太平日子。要让同一个作业扛住广告晚高峰的实时扣费与 pacing 控速、7×24 不停机,还差最后一步:它到底怎么部署上生产? 收尾这一站就把这步走完——三种部署模式怎么选、凭什么跑在 Kubernetes 上、JobManager 怎么摘掉单点,以及怎么在换版本的同时不把攒了几天的状态清零。

一句话定位:部署要回答的其实是两个正交问题——main() 在哪跑(决定隔离性)、集群交给什么资源编排后端(决定弹性与运维),而生产上的答案高度一致:Application 模式 + Flink on Kubernetes(Operator 托管)+ checkpoint / savepoint 落远端存储 + stop-with-savepoint 升级。

TL;DR

Table of contents

Open Table of contents

1. 三种部署模式:main() 在哪跑,集群是否隔离

Flink 的「部署模式」本质上只在回答两个正交的问题:① 作业的 main() 方法在哪里执行?② 集群是多作业共享,还是每个作业独占一套? 三种模式不过是这两个答案的不同组合,而它们在隔离性与运维成本上的差距,比名字听起来要大得多。

Flink 两种部署模式对比图:Session 与 Application。左侧 Session 模式(多作业共享一个常驻集群),自上而下的流程:Client(运行 main()、构建 JobGraph)→ 常驻共享集群(提前起好、长期存在)→ JobManager(1 个,共享)→ Job1/Job2/Job3(共享同一批 TaskManager/Slot);底部红色提示:一个作业 OOM 或依赖类冲突可能拖垮同集群其它作业、资源互相争抢,适合小作业、交互式 SQL、开发调试。右侧 Application 模式(每个作业独立集群,推荐生产),自上而下:Client(只上传 jar+依赖,不跑 main())→ JobManager(每作业独立,main() 在 JM 上执行,加粗高亮)→ 专属 TaskManager(随作业创建/销毁)→ 本应用唯一的 Job(资源、依赖、生命周期隔离);底部绿色提示:作业与集群同生命周期、互不影响,main() 的重活从 Client 挪到集群,避免 Client 成为提交瓶颈,生产长作业首选。

Session 省资源但作业间无隔离;Application 每作业独立集群且 main() 在 JM 上跑,隔离性最好,是生产首选。

模式main() 在哪跑集群隔离生命周期适用场景
SessionClient❌ 多作业共享集群常驻,作业来去开发调试、交互式 SQL、大量小作业
Per-Job(已弃用)Client✅ 每作业独立与作业同生命周期已被 Application 取代,了解即可
ApplicationJobManager✅ 每作业独立与作业同生命周期生产长作业首选

1.1 Session 模式:多作业共享常驻集群

Session 集群提前起好、长期存在,作业一个个提交进去,共享同一套 JobManager 和 TaskManager。它的好处很实在:集群只拉起一次,作业提交后几秒就能跑图,一堆小作业挤在一起资源利用率也高。

代价则是没有隔离——某个作业把 TaskManager 的内存吃爆触发 OOM,同集群的其它作业很可能被一起带下去,不同作业的依赖版本也会互相冲突。所以它的位置很清楚:开发调试、交互式 SQL 探索、大量短小作业,前面几篇在本地跑的正是这样一套 Session 集群。

1.2 Per-Job 模式:已弃用

Per-Job 为每个作业单独拉起一套集群,作业一结束集群随之销毁,隔离问题算是解决了;但 main() 仍然留在 Client 端执行——先在本地构建出 JobGraph 再提交,Client 一旦要同时铺开大量作业就成了瓶颈。Flink 1.15 起它已被标记为弃用(deprecated),能力被 Application 模式完整覆盖,今天知道它存在过就够了。

1.3 Application 模式:生产推荐

Application 模式在「每作业一个专属集群」之上多做了一件事:把 main() 的执行从 Client 挪到了 JobManager 上。就这一步挪动,换来了两个生产上真正在意的性质:

# Application 模式提交到 K8s(Native)
flink run-application \
  -t kubernetes-application \
  -Dkubernetes.cluster-id=ad-billing \
  -Dkubernetes.container.image=my-registry/flink-ad-billing:1.0 \
  local:///opt/flink/usrlib/ad-billing.jar

一句话选型:生产长作业一律 Application 模式;只有开发调试与交互式探索才留给 Session。

2. 部署后端:为什么是 Kubernetes

到这里定下的只是 main() 与集群的归属关系,作业最终还得踩在某个资源编排后端上跑起来,而这三个选项的取舍基本由基础设施现状决定:

后面的篇幅全给 Kubernetes:广告实时链路上云时它几乎是不需要讨论的默认答案,真正需要讨论的是上了 K8s 之后怎么跑。

Flink 跑在 K8s 上有两条成熟路线,分水岭只有一个:谁来管作业的生命周期——是你自己用命令一次次推,还是交给一个控制器持续收敛。

Flink on Kubernetes 的两条路线对比图。左侧 Native Kubernetes(Flink 直接对话 API Server),自上而下:flink run-application(-t kubernetes-application)→ Kubernetes API Server → JobManager Pod(内置 K8s ResourceManager)→ TaskManager Pods(按需自动申请/释放);底部黄色提示:Flink 自己向 K8s 申请 Pod,无需额外组件,但命令式提交,升级/HA/savepoint 编排要自己做。右侧 Kubernetes Operator(声明式 CR 管作业),自上而下:FlinkDeployment CR(一份 YAML 声明作业)→ Flink Kubernetes Operator(控制器 watch CR,加粗高亮)→ JobManager / TaskManager(由 Operator 创建管理)→ 自动 savepoint 升级与 HA(声明式、GitOps 友好);底部绿色提示:把「怎么部署/升级」交给 Operator,改 YAML 后 Operator 自动 stop-with-savepoint 再拉起,生产推荐。

Native K8s 让 Flink 自己向 API Server 要 Pod,命令式、单作业够用;Operator 用声明式 CR 托管作业的部署与升级,生产更省心。

3.1 Native Kubernetes

Native 路线靠的是 Flink 内置的一个 Kubernetes 版 ResourceManager:JobManager 直接和 K8s API Server 对话,按作业并行度的需要动态申请、释放 TaskManager Pod。一条 flink run-application -t kubernetes-application 就能把作业拉起来,不依赖任何额外组件,Flink 自己搞定 Pod 申请与弹性。

短板同样明确:提交动作是命令式的,升级、HA 编排、savepoint 管理都得你自己写脚本串成流程。作业数量不多、运维链路简单时它足够用;一旦作业上了几十个,这些脚本迟早变成运维负债。

3.2 Kubernetes Operator(生产推荐)

flink-kubernetes-operator 是官方给出的声明式答案:你把一个作业描述成一份 FlinkDeployment 自定义资源(CR) 的 YAML,Operator 作为控制器持续 watch 这个 CR,不断让集群的实际状态向 YAML 声明的目标收敛。

apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
  name: ad-billing
spec:
  image: my-registry/flink-ad-billing:1.0
  flinkVersion: v1_20
  flinkConfiguration:
    taskmanager.numberOfTaskSlots: "4"
    state.savepoints.dir: s3://flink/savepoints/ad-billing
    state.checkpoints.dir: s3://flink/checkpoints/ad-billing
    high-availability.type: kubernetes
  jobManager:
    resource: { memory: "2048m", cpu: 1 }
  taskManager:
    resource: { memory: "4096m", cpu: 2 }
  job:
    jarURI: local:///opt/flink/usrlib/ad-billing.jar
    parallelism: 8
    upgradeMode: savepoint          # 升级时自动 stop-with-savepoint 再从 savepoint 恢复

4. JobManager 高可用与 checkpoint 远端存储

Native 还是 Operator,决定的只是作业怎么被拉起;下面这两件地基工程不做,作业在本地跑得再顺,上了生产也迟早翻车。

4.1 JobManager HA:消除单点

JobManager 是整个作业的协调中枢(架构篇):调度 subtask、触发 checkpoint、汇总快照元数据都归它管,单个 JM 一挂,作业就地停摆。K8s 上的解法是 Kubernetes 高可用服务——用基于 ConfigMap 的 leader 选举记录「谁是主 JM」以及作业元数据,主 JM 失联后剩余候选自动选出新主,再从最近一次完成的 checkpoint 把作业恢复回来。配置只有两行:

high-availability.type: kubernetes
high-availability.storageDir: s3://flink/ha/ad-billing

storageDir 里存的是 JobGraph、已完成 checkpoint 的指针这类元数据,因此它必须落在持久存储上——新主 JM 要能读到它才接得住作业,写在 Pod 本地盘等于没写。

4.2 状态快照必须落远端

这就引出第二件地基:Pod 是易失的,被调度走、被重启、被驱逐,本地盘上的东西一并消失。所以 checkpoint 与 savepoint 必须写到远端持久存储(S3 / OSS / HDFS),而不是容器里的某个路径:

state.checkpoints.dir: s3://flink/checkpoints/ad-billing
state.savepoints.dir:  s3://flink/savepoints/ad-billing

大状态作业还要把状态篇的结论一并带上:RocksDB 状态后端 + 增量 checkpoint,只把本地 RocksDB 新增的那部分异步上传到远端,这样几十 GB 的归因状态既扛得住,又不会让每轮 checkpoint 的耗时反过来拖垮作业。

5. 资源与内存模型:别被 K8s OOMKill

HA 与远端存储保证的是「作业挂了能活回来」,这一节要保证的是「Pod 别先被杀掉」。容器里跑 Flink 最常见的翻车方式,就是内存参数和 Pod limit 对不齐 → 进程被 K8s OOMKill;想配平它,先得看清 Flink 的内存分层:

关键纪律只有一条:Pod 的 memory limit 必须 ≥ Total Process Memory,并且给 JVM overhead 留足余量。用 RocksDB 时尤其要留神,Managed Memory 走的是堆外,而容器统计的是 RSS(含堆外),只按 JVM Heap 去配 limit,作业跑上几小时必被 OOMKill。

并行度与 Slot 的账也要在这一步算清:source 并行度别超过 Kafka 分区数(多出来的 subtask 只会空转),再拿 TaskManager 的 numberOfTaskSlots 去除总并行度,才知道该开多少个 TM Pod——这正是架构篇里 Slot 与并行度那节在 K8s 上的落地。

6. 优雅升级与可观测

资源配平之后作业能稳定跑,但流作业的日常恰恰是不停地改:加一个计费维度、给热 key 扩并行度、换一版镜像。这时候怎么改,决定了你会不会亲手把攒了几天的状态清零。

6.1 用 savepoint 升级,不是重启

改代码、扩并行度、迁移集群,动作都是同一套:先打 savepoint(手动触发、格式稳定的状态快照,见状态篇),再从它恢复。

# 停作业并打 savepoint
flink stop --savepointPath s3://flink/savepoints/ad-billing <jobId>

# 用新版本 jar / 新并行度从 savepoint 恢复
flink run-application -t kubernetes-application \
  -Dexecution.savepoint.path=s3://flink/savepoints/ad-billing/savepoint-xxx \
  ...

用 Operator 的话,这一整套都被 upgradeMode: savepoint 接走,你只管改 YAML。流作业几乎不存在「直接重启」这个选项,因为那等于丢掉全部状态:累计消耗清零、去重记忆消失,重算出来的扣费金额没人敢认。

6.2 可观测:Metrics + UI + 日志

最后是让作业「看得见」。前面几篇排查问题时用过的那些指标,在生产上都要变成常驻的观测面:

参考


–views
Share this post on:

Previous Post
MySQL 内核深挖 · InnoDB 存储结构与 B+ 树索引
Next Post
Flink 深挖 · 反压定位、数据倾斜与性能调优