广告实时链路里,Kafka 只负责把曝光、点击、转化事件搬上流(见 Kafka 核心原理);真正要把这些事件有状态地、低延迟地、精确一次地算成钱和报表的,是下游的流处理引擎——它算重一笔,广告主当天就被多扣一笔费用;它卡住十分钟,pacing 就会在预算见底之后继续超投,大盘也停在十分钟前的数字上不动。这个系列专门深挖 Apache Flink:从运行时架构起步,依次走过时间语义与窗口、端到端 exactly-once、DataStream API、状态与 checkpoint 容错、反压定位与调优,最后落到生产部署,一路钻到能上线、能排障的深度。开篇这一篇先把地基铺平:Flink 的流处理模型长什么样,一个作业提交下去,在集群里到底是怎么跑起来的。
一句话定位:Flink 是一个以流为核心、原生有状态、支持 event-time 与 exactly-once 的分布式计算引擎;本篇讲清它的编程心智(DataStream)与运行时骨架(JobManager / TaskManager / Slot,以及从代码到物理执行的三次图变换),是读懂后面六篇的地基。
TL;DR
- 广告实时链路选 Flink 的四个硬理由:真·流(逐条处理、毫秒级延迟)、原生有状态(窗口聚合、去重、关联都靠状态)、event-time + watermark(按事件真实发生时间计算、容忍乱序迟到)、exactly-once(计费与结算不能重复也不能丢)。这四点上,微批的 Spark Streaming 和轻量的 Kafka Streams 各有短板。
- 流批一体:Flink 把批看成有界流,同一套 DataStream / Table API 既跑无界实时也跑有界回刷,广告场景常拿它做到「实时与离线对账共用一份逻辑」。
- DataStream 心智:
source → 一串 transformation(map/filter/keyBy/window/process)→ sink,数据无界、持续、逐条流动,算子则是长期驻留、不断消费上游的常驻处理单元。 - 运行时三大角色:Client(构建 JobGraph 并提交,不参与运行)、JobManager(Master,内含 Dispatcher / ResourceManager / JobMaster 三个子角色)、TaskManager(Worker,真正跑 task 的 JVM 进程)。
- Task Slot 是资源隔离单位:一个 TaskManager 划分成若干 Slot,Slot 均分该进程的托管内存(隔离内存,不隔离 CPU);作业能达到的并行度上限就是集群可用 Slot 总数。
- 代码到执行经历三次图变换:
StreamGraph(逻辑 DAG)→JobGraph(Client 侧做算子链 chaining 合并)→ExecutionGraph(JobManager 侧按并行度展开成 subtask)→ 部署到 Slot 物理执行。 - 算子链 Operator Chaining:把 forward 直连且并行度相同的相邻算子合并进同一个线程,省掉线程切换与序列化、网络开销,是 Flink 默认开启的核心优化;重算子想独占资源时可主动断链。
- Slot Sharing:默认情况下同一作业不同算子的 subtask 可以共享一个 Slot,于是 Slot 总数只需等于最大算子并行度,还天然把轻重算子搭配摆放;必要时用 slot sharing group 拆开。
- 数据交换 4 种分区策略:
forward(一对一直连)、hash(keyBy按 key 哈希路由,保证同 key 同 subtask)、rebalance(轮询均衡、治倾斜)、broadcast(每条复制到全部下游,用于广播维表与规则)。 - 部署模式:
Session(多作业共享集群)、Per-Job(已弃用)、Application Mode(推荐,main()在集群里执行、每个应用独立集群);底层可跑 on YARN 或 Kubernetes。 - AdTech 实践:实时曝光与点击统计、实时计费与 pacing、归因 join 都是 Flink 的主战场;并行度与 Slot 配置不合理(source 并行度超过 Kafka 分区数、slot sharing 把重算子挤到一起)会让资源利用严重不均,表现为部分 TaskManager 被打满、其余闲置,整体吞吐始终上不去。
Table of contents
Open Table of contents
1. 为什么广告实时链路用 Flink
广告系统的实时链路是这样一条流水线:曝光、点击、转化事件从 Kafka 流进来,下游要在几百毫秒内算清楚每个广告主花了多少钱(计费)、预算还剩多少(pacing 控速)、这次点击该归因到哪次曝光(attribution)、每个创意此刻的 CTR(实时报表)。这几件事叠在一起,对计算引擎的要求相当苛刻:延迟要压到毫秒级,因为 pacing 慢一步就是真金白银的超投;结果要准到分,因为计费错一分钱最后都要对账;还得扛得住移动端弱网补传带来的乱序与迟到。Spark Streaming 和 Kafka Streams 同样能算流,为什么广告实时链路最后普遍落到了 Flink?下面分四个维度拆开看。
1.1 真·流 vs 微批:延迟的本质差别
Spark Streaming(含 Structured Streaming 的微批模式)本质上是「攒一批再算」:它按固定间隔把流切成一个个 micro-batch,每批当作一个小批作业调度执行,因此有两个绕不开的固有代价。
- 延迟下界等于批间隔:批设 1s,端到端延迟就从 1s 起步;想再压低只能把批调小,可批越小调度开销占比越高,吞吐反而先掉下来。
- 窗口边界对不齐真实时间:微批的切分依据是处理时间,和事件自身携带的时间语义天然错位。
Flink 走的是逐条处理的真流:一条曝光事件进来就立刻穿过 source → map → keyBy → window → sink 这条算子链,中间不攒批,延迟是毫秒级,而且不会随吞吐上升而线性恶化。对 pacing 这种「预算快见底了必须立刻降速」的控制回路,毫秒与秒的差别足以决定当天会不会超投。
注:Spark 后来推出 Continuous Processing 试图补齐真流能力,但成熟度与生态都不及 Flink 的流优先设计。差异的根子在出发点——Flink 认为流是一等公民、批是流的特例,Spark 认为批是一等公民、流是批的近似,落到实时场景上体验相差很远。
1.2 原生有状态:窗口、去重、关联都靠它
延迟之外的第二道门槛是状态。广告实时计算几乎没有一处是无状态的:
- 实时曝光计数:要按 campaign 和创意分组累加,需要 keyed state 保存计数;
- 点击归因:点击事件要与前几分钟内的曝光事件做 join,需要状态缓存曝光;
- 曝光去重:同一次曝光可能被上报多次,需要状态记住已见过的 impression id;
- pacing:要实时维护每个广告主本小时已花掉多少预算,同样离不开状态。
Flink 把状态做成了引擎的一等公民:算子可以持有 keyed state 或 operator state,状态由框架托管、可落在堆内或 RocksDB、并随 checkpoint 一起容错(详见 状态管理与 Checkpoint)。相比之下,Spark Streaming 的状态支持(mapWithState、flatMapGroupsWithState)在表达力上受限,Kafka Streams 虽然也有本地状态,但绑定在 Kafka 生态内,伸缩性与状态后端的可选性都不如 Flink。
1.3 event-time 与 watermark:按事件真实时间算
有了状态,下一个问题是「这条数据该算进哪一分钟」。广告事件的到达顺序是乱序的:手机端网络抖动、批量补传,都会让一条 14:00:01 发生的点击拖到 14:00:05 才落到 Kafka。若按到达时间(processing-time)统计 14:00 这一分钟的曝光数,数字就会随集群抖动而漂移,昨天和今天重算两遍还对不上。
Flink 原生支持 event-time 语义 + watermark 机制:开窗依据事件里携带的真实发生时间,再用 watermark 度量事件时间推进到了哪里、还愿意为迟到数据多等多久。这套机制让「14:00 那一分钟究竟有多少曝光」有了确定且可复现的答案,而确定与可复现正是计费和报表的刚需(详见 时间语义、窗口与 Watermark)。
1.4 exactly-once:计费不能重复不能丢
时间语义解决了算得准不准,最后一道门槛是算得重不重。计费与结算对「每条事件被精确计算一次」是零容忍的:重复计费等于多扣广告主的钱,漏计则是平台自己少收钱。Flink 依靠 checkpoint(异步分布式快照)+ 可重放 source + 幂等或事务性 sink 这条组合拳做到端到端 exactly-once(详见 端到端 Exactly-Once)。Spark Streaming 借助幂等 sink 与 WAL 同样能给出 exactly-once,只是 Flink 的 checkpoint 机制更轻量、对端到端延迟的干扰更小。
| 维度 | Spark Streaming(微批) | Kafka Streams | Flink |
|---|---|---|---|
| 处理模型 | 微批(攒批再算) | 真流(库,嵌进应用) | 真流(独立引擎) |
| 延迟 | 秒级(≥ 批间隔) | 毫秒级 | 毫秒级 |
| 有状态 | 受限 | 支持(绑 Kafka) | 一等公民、状态后端可选 |
| event-time | 支持但弱 | 支持 | 原生 + watermark 完善 |
| exactly-once | 幂等 sink + WAL | 事务支持 | checkpoint + 2PC,端到端 |
| 部署 | 需 Spark 集群 | 嵌入式、随应用 | 独立集群 / on YARN·K8s |
结论:轻量的嵌入式实时处理选 Kafka Streams,离线为主、实时为辅选 Spark;而低延迟、有状态、event-time、exactly-once 四项全都要的广告实时链路,Flink 是当前的默认答案。
2. 流批一体与 DataStream API 心智
选型定下来之后,真正要换的是写代码的心智——批处理那套「读完、算完、结束」的思维搬到流上会处处别扭。
2.1 无界流:数据永不结束
写批处理时,你的心智是「读一个有界数据集 → 算 → 输出 → 结束」;写流处理必须换一套:数据是无界的(unbounded),持续不断地来,作业启动后常驻运行、永不主动结束。
Flink 的核心抽象是 DataStream——一条逻辑上无限的数据流。你在它上面挂一串 transformation(算子),每个算子都是长期驻留、不断消费上游并产出下游的常驻处理单元:
source(Kafka 曝光 topic)
→ map(解析成 Impression 对象)
→ filter(过滤无效流量)
→ keyBy(按 campaignId 分组)
→ window(1 分钟滚动窗口)
→ aggregate(累加曝光数)
→ sink(写 Redis / ClickHouse)
数据从 source 逐条流入,穿过这条算子链,最后由 sink 流出。这里没有跑完的概念,只有一直在跑。
2.2 流批一体:批是有界流的特例
既然流是无界的,批就可以被统一进来:Flink 把批处理看成有界流(bounded stream)——数据集有限,读到头流自然结束。于是同一套 DataStream API(以及更上层的 Table / SQL API)既能跑无界实时流,也能跑有界批,差别只在执行模式(RuntimeExecutionMode.STREAMING 与 BATCH)以及随之调整的调度与 shuffle 策略。
这个特性在广告场景里非常实用:实时计费用流模式跑,离线对账与回刷用批模式跑同一份逻辑代码,从根上避开了「实时一套 Flink、离线一套 Spark,两套代码算出两个数还得再对一遍账」这个经典痛点。
2.3 一个曝光计数的 DataStream 作业
把上面的心智落成代码,就是下面这个按 campaign 做 1 分钟滚动窗口曝光计数、再写下游的最小可运行作业:
public class ImpressionCountJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(8); // 全局默认并行度
// 1) source:从 Kafka 消费曝光事件(可重放,是 exactly-once 的前提)
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("kafka:9092")
.setTopics("ad-impressions")
.setGroupId("flink-imp-count")
.setStartingOffsets(OffsetsInitializer.committedOffsets())
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
DataStream<String> raw = env.fromSource(
source, WatermarkStrategy.noWatermarks(), "kafka-impressions");
// 2) transformation:解析 → 过滤 → 按 campaign 分组 → 开窗聚合
DataStream<ImpCount> counts = raw
.map(Impression::parse).name("parse")
.filter(imp -> imp.isValid()).name("filter-invalid")
.keyBy(Impression::campaignId) // hash 重分区:同 campaign 落同一 subtask
.window(TumblingProcessingTimeWindows.of(Time.minutes(1)))
.aggregate(new CountAgg()).name("count-1min");
// 3) sink:写入下游存储(幂等 / 事务性 sink 才能端到端 exactly-once)
counts.sinkTo(buildClickHouseSink()).name("sink-clickhouse");
env.execute("ad-impression-count"); // 提交并常驻运行
}
}
这段代码里的 map、filter、keyBy、window、aggregate、sink 各自都是一个 transformation。需要留意的是,在 env.execute() 之前什么都还没真正运行——你只是在描述一张数据流图;真正被 execute() 触发的,是把这张图编译、优化、提交到集群的整个过程(§4)。
注:
keyBy不是普通算子,它是一次重分区(redistribution),决定了哪条数据去哪个下游 subtask,既是后面 §6 分区策略的核心,也是 状态 里 keyed state 能成立的前提。
3. 运行时架构:JobManager、TaskManager 与 Slot
env.execute() 把图交出去之后,作业就进入了 Flink 的运行时。一个 Flink 集群由三类角色组成——Client、JobManager、TaskManager。
Flink 运行时三大角色:Client 提交、JobManager(Dispatcher/ResourceManager/JobMaster 三合一)调度、TaskManager 出 Slot 干活。Slot 是内存资源的隔离单位,一个作业能跑的最大并行度就等于集群可用 Slot 总数。
3.1 Client:构建并提交,不参与运行
Client 是提交作业的入口:它执行你的 main(),把 DataStream 的描述编译成 StreamGraph、再优化成 JobGraph(§4),然后把 JobGraph 提交给 JobManager。Client 本身不参与作业的实际运行——提交完毕,在 detached 模式下它就可以退出;Application Mode 更进一步,连 main() 都挪进集群里执行(见 §7)。
3.2 JobManager:集群的大脑(三个子角色)
JobManager 是 Master 进程,也是整个作业的协调中枢,内部由三个职责清晰的子角色组成:
- Dispatcher:集群的前台,接收 Client 提交的作业,为每个作业拉起一个独立的 JobMaster,并对外提供 Web UI 与 REST API。
- ResourceManager:管资源(Slot),它掌握集群里有哪些 TaskManager、各自还剩多少空闲 Slot,JobMaster 要资源时向它申请;跑在 YARN 或 K8s 上时,它还负责动态向底层资源框架申请与释放 TaskManager。
- JobMaster(旧称 JobManager 里的调度部分):管单个具体作业,负责把 ExecutionGraph 的 subtask 部署到申请到的 Slot 上、驱动并协调 checkpoint、监控 task 状态、失败时触发 restart。一个作业对应一个 JobMaster。
注:这三者的拆分正是理解部署模式(§7)的钥匙——Session 模式下一套 Dispatcher 与 ResourceManager 常驻并服务多个作业(每个作业仍各配一个 JobMaster),Application 模式下则是每个应用独立一整套。
3.3 TaskManager:真正干活的 Worker
TaskManager(TM)是 Worker 进程,每个 TM 是一个独立的 JVM,负责真正执行 subtask、缓存数据并与上下游交换数据。它启动时向 ResourceManager 注册自己的 Slot,运行中不断向 JobMaster 汇报 task 状态并发送心跳。集群的计算能力,本质上就是所有 TaskManager 的 Slot 总和。
3.4 Task Slot:资源隔离的单位
一个 TaskManager 会把自己的资源切成若干 Task Slot,理解 Slot 抓住三点即可:
- Slot 是内存的隔离单位:TM 的托管内存(managed memory)均分给各个 Slot,切 3 个 Slot 就各拿三分之一,某个 Slot 里的 task 吃不到别的 Slot 那一份,从而避免单个 task 内存膨胀拖垮整个 TM。
- Slot 不隔离 CPU:Slot 只切内存、不切 CPU,同一个 TM 里所有 Slot 的 task 共享 CPU 核,所以「一个 TM 配几个 Slot」通常参考核数(经验值是 Slot 数 ≈ CPU 核数,让每个 Slot 的 task 大致对应一个核)。
- Slot 决定并行度上限:一个作业能达到的最大并行度等于集群里所有 TM 的 Slot 总数;开启 slot sharing 时,这条约束表现为「最大算子并行度 ≤ Slot 总数」(见 §5)。Slot 不够,作业根本起不来,报
NoResourceAvailableException。
一句话记住:TaskManager 是一台干活的机器(JVM 进程),Slot 是这台机器上切出来的一个带独立内存配额的工位。 后面要反复排查的「某些 TM 被打满、某些几乎闲置」,根子多半就在 Slot 的划分与摆放没配好。
4. 从代码到执行:三次图变换
角色就位,还差一条链路:你写的 DataStream 代码,到集群里真正跑起来的物理任务之间,隔着三次图变换。看懂这三张图,才算真正理解「我这几行代码最后变成了几个线程、分别跑在哪」。
三次图变换:StreamGraph(逻辑)→ JobGraph(Client 侧做算子链合并)→ ExecutionGraph(JobManager 侧按并行度展开成 subtask)→ 部署到 Slot 物理执行。并行度决定每个算子展开成几个 subtask。
4.1 StreamGraph:逻辑 DAG
main() 里每一次 transformation 调用,都会在 StreamExecutionEnvironment 中登记一个节点;execute() 触发时,Flink 先把这些节点连成 StreamGraph——一张最原始的逻辑有向无环图,一个算子对应一个节点,忠实反映你写下的拓扑。
4.2 JobGraph:算子链合并(Client 侧)
StreamGraph 还没有经过任何优化,Client 会把它转成 JobGraph,其中最核心的动作就是算子链 Operator Chaining(§5):把满足条件的相邻算子(forward 直连、并行度相同等)合并成一个 JobVertex,既减少节点数,也省掉算子之间的序列化与线程切换。
上图里 Source 与 map 被合并成链节点 [Source → map],而 window(前面隔着 keyBy 重分区)和 sink 各自独立。JobGraph 就是 Client 最终提交给 JobManager 的东西。
4.3 ExecutionGraph:按并行度展开成 subtask(JobManager 侧)
JobManager 收到 JobGraph 后,把它展开成 ExecutionGraph,也就是物理执行图,关键动作是按并行度(parallelism)把每个 JobVertex 展开成多个并行的 subtask:
- 一个算子的并行度是几,它就被展开成几个 subtask,每个 subtask 是一个独立执行实例,处理数据的一个分片,跑在某个 Slot 内的一个线程里。
- 上图中
[Source→map]并行度 2 展开出 2 个 subtask,window并行度 2 展开出 2 个 subtask,sink并行度 1 只有 1 个 subtask。 - subtask 之间的数据传输由分区策略决定(§6):
[Source→map]到window之间隔着keyBy,是 hash 全连接重分区(all-to-all);window到sink则是收敛。
三张图一句话概括:StreamGraph 是你写了什么,JobGraph 是合并优化后提交了什么,ExecutionGraph 是集群里实际并行跑成了什么。 并行度正是从 JobGraph 放大到 ExecutionGraph 这一步的关键旋钮,它直接决定你需要多少 Slot、每个算子摊到几个线程。
5. 算子链与 Slot 共享
上一节里一带而过的 chaining,值得单独展开——它和 Slot 共享一起,决定了同样一份 ExecutionGraph 最终会以怎样的密度落在集群上。
5.1 Operator Chaining:为什么要合并算子
算子链是 Flink 默认开启的核心性能优化:把满足条件的相邻算子塞进同一个线程里串行执行,一条数据在链内算子之间直接以方法调用传递,而不必走「序列化 → 网络或内存队列 → 反序列化」这一圈。
合并带来三重收益:
- 省线程切换:链内算子共用一个线程,没有跨线程调度开销;
- 省序列化与数据交换:链内传递的是对象引用,零序列化、零网络;
- 更好的 cache 局部性:一条数据被一串算子连续处理完,数据还留在 CPU cache 里。
合并的条件(简化说)是:相邻算子并行度相同、连接方式为 forward(一对一)、处在同一个 slot sharing group,且 chaining 没有被显式关掉。上图里的 Source → map 就是典型——并行度同为 2、forward 直连,于是合并成一条链;而 map → window 之间夹着 keyBy(hash 重分区而非 forward),链在这里断开。
5.2 什么时候要主动断链
默认合并总体是好事,但有两种场景值得主动断链:
- 重算子需要独占线程、也需要更清晰的资源核算:某个
process算子如果逻辑很重(复杂计算或访问外部系统),把它和别人链在一起会拖慢整条链,也很难单独调优,这时用disableChaining()让它独立成 task。 - 想在 Web UI 里单独观测某个算子的负载与反压:链合并之后 UI 上看到的是整条链的指标,粒度太粗;断链后才能单独看这个算子的忙碌度与反压情况。
DataStream<Enriched> out = stream
.map(new LightParse()) // 轻算子
.map(new HeavyEnrich()) // 重算子:查外部维表、CPU 密集
.disableChaining() // 断链:让它独立成 task、独占线程、单独可观测
.name("heavy-enrich");
// 也可以用 startNewChain() 从某算子开始一条新链,或
// env.disableOperatorChaining() 全局关闭(一般只用于调试)
5.3 Slot Sharing:让 Slot 数只需等于最大并行度
chaining 解决的是「算子怎么合并进线程」,Slot 共享解决的则是「这些 task 怎么摆进 Slot」:同一个作业里不同算子的 subtask,默认可以共享同一个 Slot(它们默认都在名为 default 的 slot sharing group 里)。
这带来两个好处:
- Slot 总数只需等于最大算子并行度,而不是所有算子并行度之和。假设作业里 source 并行度 2、window 并行度 8、sink 并行度 2,没有 slot sharing 就要 2+8+2=12 个 Slot;有了 slot sharing,一个 Slot 里可以同时放下 source、window、sink 各一个 subtask,只需 8 个 Slot(即最大并行度)。
- 轻重算子自动搭配:一个 Slot 里既有轻量的 source subtask,又有吃资源的 window subtask,整体资源利用更均衡,不至于出现轻算子的 Slot 空转、重算子的 Slot 过载。
代价与陷阱同样存在:slot sharing 也可能把几个重算子的 subtask 一起塞进同一个 Slot,导致该 Slot 所在的 TM 被打满、别的 TM 反倒相对空闲。这时的解法是用 slot sharing group 把某些算子隔离到独立的 Slot 组:
stream
.map(new HeavyCpuOp()).slotSharingGroup("heavy") // 重算子单独一个 slot 组
.map(new LightOp()); // 回到默认组(也可显式 .slotSharingGroup("default"))
注:划到不同 slot sharing group 的算子,其 subtask 不会共享 Slot,相当于给重算子单开一间包间。这既是隔离资源的常规手段,也是修复 §3.4 那种「部分 TM 打满、其余闲置」的关键动作之一。
6. 数据交换与分区策略
算子怎么摆放讲完了,还剩最后一个物理层面的问题:subtask 之间到底怎么传数据。这由**分区策略(partitioning)**决定,也就是「上游某个 subtask 产出的一条数据,该发给下游哪一个或哪几个 subtask」,它直接影响数据分布、状态正确性以及会不会倾斜。最常用的是下面四种。
四种分区策略:forward 一对一(算子链前提)、hash/keyBy 按 key 哈希(同 key 同 subtask,keyed state 的基础)、rebalance 轮询均衡(治倾斜)、broadcast 广播(每条复制到全部,用于广播维表/规则)。
- forward(一对一):上游第
i个 subtask 直连下游第i个 subtask,要求上下游并行度相同,开销最小,也是 §5.1 里算子链合并的前提;相邻的同并行度算子默认就是 forward。 - hash /
keyBy:按hash(key) % parallelism决定数据去哪个下游 subtask,核心保证是相同 key 的所有数据一定落到同一个下游 subtask——这正是 keyed state(按 key 维护状态)以及按 key 聚合、去重、开窗能正确工作的根本前提。§2.3 里的keyBy(campaignId)保证了同一个 campaign 的曝光永远进同一个 subtask 累加,计数才不会算错。 - rebalance(轮询):上游以 round-robin 把数据均匀发给所有下游 subtask,用来打散数据倾斜——当某个 key 特别热导致单个 subtask 过载,而下游算子又不需要按 key 分组时,
rebalance()能把负载摊平;同族的rescale只在本地组内轮询,开销更小。 - broadcast(广播):上游每条数据复制发送到所有下游 subtask,用于广播小维表、规则与配置——例如把计费规则表广播给所有做归因的 subtask,让每个 subtask 本地都持有一份完整规则(Broadcast State,详见 状态管理)。
注:回到广告链路,曝光计数必须用
keyBy(campaignId)(hash,同 campaign 同 subtask 才能正确累加);若少数超大 campaign 造成倾斜、而某一步又不需要按 key,就用rebalance把负载摊平;广告主预算阈值这类小表则用broadcast发给所有 subtask,供本地 pacing 判断。选错分区策略的后果只有两种:该 keyBy 却 rebalance 会算错,该 rebalance 却放任热 key 会压垮单个 subtask。
7. 部署模式:Session、Per-Job 与 Application Mode
同一个作业、同一张 ExecutionGraph,还可以用不同的部署模式跑起来,区别在于集群的生命周期与**main() 在哪里执行**,而这两点又直接决定了资源隔离粒度和日常运维方式。
- Session 模式:先起一个长期运行的集群(一套 Dispatcher 与 ResourceManager 常驻),再往里提交多个作业,多作业共享同一批 TaskManager。好处是作业启动快、资源复用充分;代价是隔离性差——一个作业把 TM 搞挂(例如 OOM)会连累同集群的其他作业,资源竞争也互相影响。适合大量短小作业,或开发测试环境。
- Per-Job 模式(已弃用):每个作业独占一个集群,作业结束即销毁集群,隔离性不错,但
main()仍在 Client 侧执行、Client 负担很重。Flink 1.15 起已标记弃用,由 Application Mode 接棒。 - Application 模式(推荐):每个应用独立一套集群(强隔离),关键区别在于
main()挪到集群里(JobManager 上)执行而不再留在 Client 侧。这样 Client 极轻,只管把 jar 与依赖丢过去,生成 JobGraph 的开销、下载依赖的开销全在集群内完成,特别适合生产。
底层资源框架上,以上模式都可以跑在不同的资源管理器之上:
- on YARN:Flink 的 ResourceManager 向 YARN 申请 container 作为 TaskManager,是 Hadoop 生态的常规选择。
- on Kubernetes:Flink 原生支持 K8s(Native Kubernetes),ResourceManager 直接向 K8s API Server 申请 Pod 作为 TaskManager,弹性伸缩,是当前云原生环境的主流。
- Standalone:手动部署一批 JM / TM 进程,不依赖外部资源框架,简单但缺乏弹性。
生产建议:核心实时作业用 Application Mode + on K8s——强隔离(一个作业崩了不影响别人)、
main()在集群内执行、还能弹性伸缩;Session 模式留给那些数量多、体量小、能容忍互相影响的作业。
8. 一个作业从提交到运行的完整流程
最后把前七节的角色、图变换、分区策略与部署模式串成一条线,完整走一遍从 env.execute() 到数据开始流动的生命周期:
- Client 编译:执行
main(),先生成 StreamGraph,再做算子链合并得到 JobGraph(Application 模式下这一步在集群内完成)。 - 提交作业:Client 把 JobGraph 连同用户 jar 与依赖一起提交给 Dispatcher。
- 拉起 JobMaster:Dispatcher 为这个作业启动一个专属 JobMaster,并把 JobGraph 交给它。
- 展开 ExecutionGraph:JobMaster 按各算子并行度把 JobGraph 展开成 ExecutionGraph(一堆 subtask),同时算出需要多少 Slot。
- 申请资源:JobMaster 向 ResourceManager 申请 Slot;跑在 YARN 或 K8s 上时,Slot 不够则由 ResourceManager 动态拉起新的 TaskManager。
- 分配 Slot:ResourceManager 在考虑 slot sharing 的前提下,把空闲 Slot 分配给 JobMaster。
- 部署 subtask:JobMaster 把各个 subtask 部署到对应 TaskManager 的 Slot 上,TM 随即启动 task 线程。
- 建立数据通道并运行:subtask 之间按分区策略(forward / hash / rebalance / broadcast)建立数据传输通道,source 开始从 Kafka 拉数据,整条流水线持续处理。
- 常驻运行与容错:运行期间 JobMaster 周期性触发 checkpoint 并监控 task 心跳,task 失败则触发 restart、从最近一次 checkpoint 恢复状态(详见 状态与 Checkpoint)。
注:这条链路上任何一环卡住,都会如实反映到 Web UI 上——卡在第 6 步报
NoResourceAvailableException就是 Slot 不够;跑起来之后部分 subtask 忙、部分闲,则是并行度、分区策略或 slot sharing 没配好。地基到这里就铺完了,下一篇会回到 §2.3 那个被一笔带过的window,讲清楚事件时间究竟怎么推进、迟到数据还能等多久。
参考
- Apache Flink. Flink Architecture(Flink 官方 · 运行时架构):JobManager/TaskManager/Slot、部署模式的一手权威说明。
- Apache Flink. Flink DataStream API Programming Guide:DataStream 编程模型、transformation、source/sink 官方指南。
- Apache Flink. Jobs and Scheduling(StreamGraph→JobGraph→ExecutionGraph 与调度):三次图变换与 subtask 调度的内部机制。
- Apache Flink. Flink and Deployment Modes(Session / Application Mode):Session / Application 部署模式与 on YARN·K8s 的官方对比。
- Fabian Hueske, Vasiliki Kalavri. Stream Processing with Apache Flink(O’Reilly):Flink 流处理模型与运行时的系统性经典读物。