前面六站——架构、时间语义与窗口、端到端 Exactly-Once、DataStream、状态与容错与反压调优——都在本地的 docker / Session 集群里把作业调到了「跑得对、跑得快」,可那毕竟是一台笔记本上的太平日子。要让同一个作业扛住广告晚高峰的实时扣费与 pacing 控速、7×24 不停机,还差最后一步:它到底怎么部署上生产? 收尾这一站就把这步走完——三种部署模式怎么选、凭什么跑在 Kubernetes 上、JobManager 怎么摘掉单点,以及怎么在换版本的同时不把攒了几天的状态清零。
一句话定位:部署要回答的其实是两个正交问题——
main()在哪跑(决定隔离性)、集群交给什么资源编排后端(决定弹性与运维),而生产上的答案高度一致:Application 模式 + Flink on Kubernetes(Operator 托管)+ checkpoint / savepoint 落远端存储 + stop-with-savepoint 升级。
TL;DR
- 三种部署模式,差别只在两件事:
main()在哪跑、集群是否独占——Session(多作业共享常驻集群)、Per-Job(每作业独立集群,已弃用)、Application(每作业独立集群且main()在 JobManager 上跑,生产推荐)。 - Session 模式:集群提前起好、多个作业共享一套 JobManager 与 TaskManager,省资源、启动快,代价是一个作业 OOM 或异常就可能拖垮同集群的其它作业,资源还互相争抢——只适合小作业、交互式 SQL 与开发调试。
- Application 模式:每个应用一个专属集群,且
main()从 Client 挪到 JobManager 上执行,于是作业与集群同生命周期、资源与依赖彻底隔离,Client 也不再是提交瓶颈——生产长作业首选。 - 部署后端:底层可跑在
Standalone/YARN/Kubernetes上,新项目基本都落在 Kubernetes。 - Flink on K8s 两条路线:Native Kubernetes(Flink 自己对话 API Server、按需申请 Pod)与 Kubernetes Operator(
flink-kubernetes-operator,用声明式FlinkDeploymentCR 托管作业);生产推荐 Operator,声明式、GitOps 友好,savepoint 升级与 HA 都被它接走。 - JobManager 高可用:单 JM 就是单点,K8s 上用 Kubernetes HA 服务(基于 ConfigMap 的 leader 选举加元数据存储),JM 挂掉后自动选主,并配合 checkpoint 把作业恢复回来。
- checkpoint / savepoint 必须落远端:Pod 是易失的,状态快照要写到 S3 / OSS / HDFS 这类持久存储,本地盘不能当归宿。
- 内存模型要和容器配对:TaskManager 内存自外向内分
Total Process→Total Flink→Managed(RocksDB、批算法用)等层次,容器里跑 Flink 必须把 Pod limit 和 Flink 内存参数对齐,否则直接被 K8s OOMKill。 - 升级靠 savepoint:改代码、扩并行度都走
stop-with-savepoint→ 从 savepoint 重启(承接状态篇),Operator 能把这套流程自动化。 - 可观测:metrics 推给 Prometheus 再用 Grafana 建大盘、Flink Web UI 盯反压面板(承接反压篇)、容器日志集中收集。
Table of contents
Open Table of contents
1. 三种部署模式:main() 在哪跑,集群是否隔离
Flink 的「部署模式」本质上只在回答两个正交的问题:① 作业的 main() 方法在哪里执行?② 集群是多作业共享,还是每个作业独占一套? 三种模式不过是这两个答案的不同组合,而它们在隔离性与运维成本上的差距,比名字听起来要大得多。
Session 省资源但作业间无隔离;Application 每作业独立集群且 main() 在 JM 上跑,隔离性最好,是生产首选。
| 模式 | main() 在哪跑 | 集群隔离 | 生命周期 | 适用场景 |
|---|---|---|---|---|
| Session | Client | ❌ 多作业共享 | 集群常驻,作业来去 | 开发调试、交互式 SQL、大量小作业 |
| Per-Job(已弃用) | Client | ✅ 每作业独立 | 与作业同生命周期 | 已被 Application 取代,了解即可 |
| Application | JobManager | ✅ 每作业独立 | 与作业同生命周期 | 生产长作业首选 |
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 上。就这一步挪动,换来了两个生产上真正在意的性质:
- 隔离彻底:集群与作业同生命周期,一个作业的资源、依赖与故障都关在自己的集群里,晚高峰某个作业把 Slot 打满,也波及不到隔壁的归因 join 作业。
- Client 变轻:Client 只负责上传 jar 与依赖、拉起集群,不必再执行可能很重的
main()(有些作业的main()里要拉配置、下载维表、拼装大量算子)。批量提交时,Client 不再是瓶颈。
# 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() 与集群的归属关系,作业最终还得踩在某个资源编排后端上跑起来,而这三个选项的取舍基本由基础设施现状决定:
- Standalone:在物理机或 VM 上手动起 JM 与 TM 进程,胜在简单,但弹性伸缩与故障恢复全靠自己兜,生产少用。
- YARN:Hadoop 生态里的经典答案,和离线集群共用一个资源池,存量 Hadoop 环境里仍然常见。
- Kubernetes:新项目的主流——容器编排、弹性伸缩、故障自愈、声明式运维都是现成的,还能和公司其它微服务共用同一套基础设施(承接容器化与 K8s 部署)。
后面的篇幅全给 Kubernetes:广告实时链路上云时它几乎是不需要讨论的默认答案,真正需要讨论的是上了 K8s 之后怎么跑。
3. Flink on Kubernetes:Native 与 Operator 两条路
Flink 跑在 K8s 上有两条成熟路线,分水岭只有一个:谁来管作业的生命周期——是你自己用命令一次次推,还是交给一个控制器持续收敛。
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 恢复
- 声明式 / GitOps 友好:作业定义进 Git,改 YAML 就是改作业,天然可评审、可审计、可回滚。
- 自动化升级与 HA:
upgradeMode: savepoint让 Operator 在你改镜像或并行度时自动stop-with-savepoint→ 用新配置从 savepoint 拉起,把状态篇里手动敲的那套 savepoint 升级流程(§6.1 会再走一遍)整个接管。 - 生产推荐:作业一多,声明式托管省下的运维成本远超 Native 那条命令带来的便利。
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 的内存分层:
- Total Process Memory(进程总内存,要 ≤ Pod 的 memory limit)
- Total Flink Memory
- JVM Heap(算子、用户代码、堆内状态后端)
- Managed Memory(RocksDB、批算法、排序用的堆外内存)
- Network Buffers(算子间数据交换的缓冲,反压篇里 credit-based 流控反复申请与释放的正是这块)
- JVM Overhead / Metaspace(JVM 自身开销)
- Total Flink Memory
关键纪律只有一条: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 + 日志
最后是让作业「看得见」。前面几篇排查问题时用过的那些指标,在生产上都要变成常驻的观测面:
- Metrics:Flink 内置 metrics(吞吐、延迟、checkpoint 大小与耗时、反压指标
busyTime/backPressuredTime)暴露给 Prometheus,再用 Grafana 建大盘。 - Web UI:算子链、并行度、checkpoint 历史与反压面板(High/OK)都在这里,是线上定位问题的第一现场(反压篇)。
- 日志:容器日志集中收集(如 Loki / ELK),配合 traceId 关联上下游。
参考
- Apache Flink. Deployment Modes(Application / Session):三种部署模式的官方定义与取舍。
- Apache Flink. Native Kubernetes:Flink 直接对接 K8s API Server 的部署方式。
- Apache Flink. Kubernetes Operator(flink-kubernetes-operator):声明式
FlinkDeploymentCR、升级模式与 HA。 - Apache Flink. High Availability (HA):JobManager 高可用与 Kubernetes HA 服务。
- Apache Flink. TaskManager Memory Model:TaskManager 内存分层与容器化配置,避免 OOMKill。