前面几站——架构、SQL 与时间窗口、真实数据源与 Exactly-Once、DataStream、状态容错与反压调优——都在本地的 docker / session 集群里”跑得对、跑得快”。但一个流作业要真正扛住广告晚高峰、7×24 不间断,就得回答最后一个问题:它到底怎么部署上生产? 这一站(也是整个系列的收尾)就讲部署:三种部署模式怎么选、怎么跑在 Kubernetes 上、JobManager 怎么高可用、怎么优雅升级不丢状态。
一句话定位:部署的核心是回答两个问题——
main()在哪跑(决定隔离性)、集群跑在什么资源编排上(决定弹性与运维)。生产答案高度一致:Application 模式 + Flink on Kubernetes(Operator)+ 远端 checkpoint 存储 + 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)vs Kubernetes Operator(
flink-kubernetes-operator,用声明式FlinkDeploymentCR 管作业)。生产推荐 Operator:声明式、GitOps 友好、自动处理 savepoint 升级与 HA。 - JobManager 高可用:单 JM 是单点。K8s 上用 Kubernetes HA 服务(基于 ConfigMap 的 leader 选举 + 元数据存储),JM 挂了自动选主恢复,配合 checkpoint 恢复作业。
- checkpoint / savepoint 必须落远端:状态快照要存到 S3 / OSS / HDFS 等持久存储,Pod 是易失的,本地盘不能作为 checkpoint 归宿。
- 内存模型要配对: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 模式:多作业共享常驻集群
集群提前起好、长期存在,然后往里提交多个作业,它们共享同一套 JobManager 和 TaskManager。
- 优点:集群只起一次,作业启动快、资源利用率高(多个小作业挤在一起)。
- 缺点:没有隔离——一个作业把 TaskManager 的内存吃爆(OOM),可能连累同集群的其它作业一起挂;不同作业的依赖版本也可能冲突。
- 适合:开发调试、交互式 SQL 探索、大量短小作业。前面几篇在本地跑的就是 Session 集群。
1.2 Per-Job 模式:已弃用
每个作业单独拉起一个集群,作业结束集群销毁。它解决了隔离问题,但 main() 仍在 Client 端执行(构建 JobGraph 再提交),Client 容易成为瓶颈。Flink 1.15 起已弃用(deprecated),被 Application 模式取代,了解即可。
1.3 Application 模式:生产推荐
每个应用一个专属集群,并且把 main() 的执行从 Client 挪到了 JobManager 上。这一步挪动带来两个关键好处:
- 隔离彻底:集群与作业同生命周期,一个作业的资源、依赖、故障都关在自己的集群里,不影响别人。
- 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
上面说的是”部署模式”,它们底层还要跑在某个资源编排后端上:
- Standalone:手动在物理机/VM 上起 JM、TM 进程。简单,但弹性、故障恢复全靠自己,生产少用。
- YARN:Hadoop 生态里的经典选择,和大数据集群共用资源池。存量 Hadoop 环境仍常见。
- Kubernetes:新项目的主流。天然的容器编排、弹性伸缩、故障自愈、声明式运维,和公司其它微服务共用一套基础设施(承接容器化与 K8s 部署)。
后面重点讲 Kubernetes,因为它是当下广告实时链路上云的默认答案。
3. Flink on Kubernetes:Native 与 Operator 两条路
Flink 跑在 K8s 上有两条成熟路线,区别在于**“谁来管作业的生命周期”**。
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 就能拉起。
- 优点:不依赖额外组件,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 升级流程全自动化了。 - 生产推荐:作业一多,声明式托管带来的运维收益远超 Native 的”一条命令”。
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 的内存分层:
- Total Process Memory(进程总内存,要 ≤ Pod 的 memory limit)
- Total Flink Memory
- JVM Heap(算子、用户代码、堆内状态后端)
- Managed Memory(RocksDB、批算法、排序用的堆外内存)
- Network Buffers(算子间数据交换缓冲,反压相关,见反压篇)
- 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 配合总并行度算好 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 + 日志
- 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。