广告实时链路里,Kafka 只是把事件搬到了流上(见 Kafka 核心原理);真正把这些曝光、点击、转化事件有状态地、低延迟地、精确一次地算成钱和报表的,是流处理引擎。这个系列专门深挖 Apache Flink——从运行时架构、时间与窗口、端到端 Exactly-Once、DataStream,到状态容错、反压调优与生产部署,一路钻到能上生产、能定位问题的深度。开篇这一篇先把地基铺平: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(算子链 chaining 合并算子、Client 侧生成)→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 慢一步就超投)、结果要准(计费不能算错一分钱)、还要能扛住乱序和迟到的事件。为什么最后大家选了 Flink,而不是 Spark Streaming 或 Kafka Streams?看四个维度。
1.1 真·流 vs 微批:延迟的本质差别
Spark Streaming(含 Structured Streaming 的微批模式)本质是”攒一批再算”:它把流切成一个个小批次(micro-batch,比如每 1 秒一批),每批当成一个小的批作业跑。这带来两个固有代价:
- 延迟下界 = 批间隔:批设 1s,端到端延迟就至少 1s 起步;想更低就得把批调小,但批越小调度开销占比越高,吞吐反而掉。
- 窗口对不齐真实时间:微批的边界是”处理时间”的批切分,和事件本身的时间语义天然错位。
Flink 是逐条处理的真流:一条曝光事件进来,立刻流过 source → map → keyBy → window → sink 这条算子链,中间不攒批。延迟是毫秒级,且延迟不随吞吐线性恶化。对 pacing 这种”预算快花完了要立刻降速”的场景,毫秒 vs 秒级是决定性的。
注: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 才到达。如果按”到达时间(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 也能做到 exactly-once(依赖幂等 sink + WAL),但 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 vs 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 要资源时向它申请,它负责分配;在 on 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 就各拿 1/3,一个 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 共享
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 数只需等于最大并行度
Flink 默认还有一个巧妙机制:Slot 共享(Slot Sharing)——同一个作业里、不同算子的 subtask,可以共享同一个 Slot(默认都在同一个 slot sharing group default)。
它带来两个好处:
- 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,等于给重算子”包间”。这既是隔离资源的手段,也是修复”重算子挤一起”的关键动作之一。
6. 数据交换与分区策略
subtask 之间怎么传数据,由**分区策略(partitioning)**决定——即”上游某个 subtask 产出的一条数据,应该发给下游哪个(些)subtask”。这直接影响数据分布、状态正确性和是否倾斜。四种最常用的:
四种分区策略:forward 一对一(算子链前提)、hash/keyBy 按 key 哈希(同 key 同 subtask,keyed state 的基础)、rebalance 轮询均衡(治倾斜)、broadcast 广播(每条复制到全部,用于广播维表/规则)。
- forward(一对一):上游第
i个 subtask 直连下游第i个 subtask,要求上下游并行度相同。开销最小、是算子链合并的前提。默认相邻同并行度算子就是 forward。 - hash /
keyBy:按 key 的哈希值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,详见 状态管理)。
AdTech 例子:曝光计数用
keyBy(campaignId)(hash,保证同 campaign 同 subtask 才能正确累加);若发现少数超大 campaign 造成倾斜、且某一步不需要按 key,就用rebalance摊平;把”广告主预算阈值配置”这类小表用broadcast发给所有 subtask 做本地 pacing 判断。选错分区策略要么算错(该 keyBy 却 rebalance),要么倾斜(该 rebalance 却让热 key 压垮一个 subtask)。
7. 部署模式:Session、Per-Job 与 Application Mode
同样一个作业,可以用不同的部署模式跑,区别在于集群的生命周期和**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(),把 DataStream 描述先生成 StreamGraph,再做算子链合并生成 JobGraph(Application 模式下这步在集群内做)。 - 提交作业:Client 把 JobGraph(+ 用户 jar、依赖)提交给 Dispatcher。
- 拉起 JobMaster:Dispatcher 为这个作业启动一个专属 JobMaster,把 JobGraph 交给它。
- 展开 ExecutionGraph:JobMaster 把 JobGraph 按各算子并行度展开成 ExecutionGraph(一堆 subtask),并计算需要多少 Slot。
- 申请资源:JobMaster 向 ResourceManager 申请 Slot;on YARN/K8s 时,ResourceManager 若 Slot 不够会动态拉起新的 TaskManager。
- 分配 Slot:ResourceManager 把空闲 Slot(考虑 slot sharing)分配给 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 配置问题。
参考
- 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 流处理模型与运行时的系统性经典读物。