Skip to content
Charles Shao
Go back

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

views

前面几站——架构SQL 与时间窗口真实数据源与 Exactly-OnceDataStream状态容错反压调优——都在本地的 docker / session 集群里”跑得对、跑得快”。但一个流作业要真正扛住广告晚高峰、7×24 不间断,就得回答最后一个问题:它到底怎么部署上生产? 这一站(也是整个系列的收尾)就讲部署:三种部署模式怎么选、怎么跑在 Kubernetes 上、JobManager 怎么高可用、怎么优雅升级不丢状态。

一句话定位:部署的核心是回答两个问题——main() 在哪跑(决定隔离性)、集群跑在什么资源编排上(决定弹性与运维)。生产答案高度一致:Application 模式 + Flink on Kubernetes(Operator)+ 远端 checkpoint 存储 + 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 模式:多作业共享常驻集群

集群提前起好、长期存在,然后往里提交多个作业,它们共享同一套 JobManager 和 TaskManager。

1.2 Per-Job 模式:已弃用

每个作业单独拉起一个集群,作业结束集群销毁。它解决了隔离问题,但 main() 仍在 Client 端执行(构建 JobGraph 再提交),Client 容易成为瓶颈。Flink 1.15 起已弃用(deprecated),被 Application 模式取代,了解即可。

1.3 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

上面说的是”部署模式”,它们底层还要跑在某个资源编排后端上:

后面重点讲 Kubernetes,因为它是当下广告实时链路上云的默认答案。

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

Flink 内置了一个 Kubernetes 版的 ResourceManager,JobManager 直接和 K8s API Server 对话,按作业需要动态申请/释放 TaskManager Pod。你只要一条 flink run-application -t kubernetes-application 就能拉起。

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 存储

生产部署有两个”不做就会翻车”的地基。

4.1 JobManager HA:消除单点

JobManager 是作业协调中枢(架构篇),单个 JM 挂了作业就停。K8s 上用 Kubernetes 高可用服务:基于 ConfigMap 的 leader 选举存储”谁是主 JM”和作业元数据,JM 挂掉后自动选出新主、从最近 checkpoint 恢复作业。配置很简单:

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

storageDir 存的是 JobGraph、已完成 checkpoint 的指针等元数据——它必须落在持久存储上,这样新主 JM 才能读到并接管。

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 的增量部分异步上传到远端,既扛得住大状态又不拖垮 checkpoint。

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

容器里跑 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 配合总并行度算好 Pod 数量——这也是架构篇 Slot 那节的落地。

6. 优雅升级与可观测

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 深挖 · 反压定位与性能调优