Skip to content
Charles Shao
Go back

Flink 深挖(开篇)· 流处理模型与运行时架构

–views

广告实时链路里,Kafka 只负责把曝光、点击、转化事件搬上流(见 Kafka 核心原理);真正要把这些事件有状态地、低延迟地、精确一次地算成钱和报表的,是下游的流处理引擎——它算重一笔,广告主当天就被多扣一笔费用;它卡住十分钟,pacing 就会在预算见底之后继续超投,大盘也停在十分钟前的数字上不动。这个系列专门深挖 Apache Flink:从运行时架构起步,依次走过时间语义与窗口、端到端 exactly-once、DataStream API、状态与 checkpoint 容错、反压定位与调优,最后落到生产部署,一路钻到能上线、能排障的深度。开篇这一篇先把地基铺平:Flink 的流处理模型长什么样,一个作业提交下去,在集群里到底是怎么跑起来的。

一句话定位:Flink 是一个以流为核心、原生有状态、支持 event-time 与 exactly-once 的分布式计算引擎;本篇讲清它的编程心智(DataStream)与运行时骨架(JobManager / TaskManager / Slot,以及从代码到物理执行的三次图变换),是读懂后面六篇的地基。

TL;DR

Table of contents

Open Table of contents

广告系统的实时链路是这样一条流水线:曝光、点击、转化事件从 Kafka 流进来,下游要在几百毫秒内算清楚每个广告主花了多少钱(计费)、预算还剩多少(pacing 控速)、这次点击该归因到哪次曝光(attribution)、每个创意此刻的 CTR(实时报表)。这几件事叠在一起,对计算引擎的要求相当苛刻:延迟要压到毫秒级,因为 pacing 慢一步就是真金白银的超投;结果要准到分,因为计费错一分钱最后都要对账;还得扛得住移动端弱网补传带来的乱序与迟到。Spark Streaming 和 Kafka Streams 同样能算流,为什么广告实时链路最后普遍落到了 Flink?下面分四个维度拆开看。

1.1 真·流 vs 微批:延迟的本质差别

Spark Streaming(含 Structured Streaming 的微批模式)本质上是「攒一批再算」:它按固定间隔把流切成一个个 micro-batch,每批当作一个小批作业调度执行,因此有两个绕不开的固有代价。

Flink 走的是逐条处理的真流:一条曝光事件进来就立刻穿过 source → map → keyBy → window → sink 这条算子链,中间不攒批,延迟是毫秒级,而且不会随吞吐上升而线性恶化。对 pacing 这种「预算快见底了必须立刻降速」的控制回路,毫秒与秒的差别足以决定当天会不会超投。

注:Spark 后来推出 Continuous Processing 试图补齐真流能力,但成熟度与生态都不及 Flink 的流优先设计。差异的根子在出发点——Flink 认为流是一等公民、批是流的特例,Spark 认为批是一等公民、流是批的近似,落到实时场景上体验相差很远。

1.2 原生有状态:窗口、去重、关联都靠它

延迟之外的第二道门槛是状态。广告实时计算几乎没有一处是无状态的:

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 StreamsFlink
处理模型微批(攒批再算)真流(库,嵌进应用)真流(独立引擎)
延迟秒级(≥ 批间隔)毫秒级毫秒级
有状态受限支持(绑 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 客户端负责构建 JobGraph 并提交作业,它本身不是运行时组件(提交完即可退出,Application Mode 下 main 甚至在集群内运行)。中间紫色面板是 JobManager(Master 进程),内含三个子角色:Dispatcher 负责接收作业、为每个作业拉起一个 JobMaster 并提供 Web UI;ResourceManager 负责管理和申请 Task Slot、对接 YARN/Kubernetes 等资源框架;JobMaster 负责调度单个具体作业、把 subtask 部署到 slot、协调 checkpoint。下方两个 TaskManager 面板是 Worker 工作进程,每个是一个独立 JVM,内部划分为 3 个 Task Slot(资源隔离单位)。JobManager 向下用箭头连到两个 TaskManager,标注为部署 task 与心跳/申请 slot。

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 进程,也是整个作业的协调中枢,内部由三个职责清晰的子角色组成:

注:这三者的拆分正是理解部署模式(§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 抓住三点即可:

  1. Slot 是内存的隔离单位:TM 的托管内存(managed memory)均分给各个 Slot,切 3 个 Slot 就各拿三分之一,某个 Slot 里的 task 吃不到别的 Slot 那一份,从而避免单个 task 内存膨胀拖垮整个 TM。
  2. Slot 不隔离 CPU:Slot 只切内存、不切 CPU,同一个 TM 里所有 Slot 的 task 共享 CPU 核,所以「一个 TM 配几个 Slot」通常参考核数(经验值是 Slot 数 ≈ CPU 核数,让每个 Slot 的 task 大致对应一个核)。
  3. Slot 决定并行度上限:一个作业能达到的最大并行度等于集群里所有 TM 的 Slot 总数;开启 slot sharing 时,这条约束表现为「最大算子并行度 ≤ Slot 总数」(见 §5)。Slot 不够,作业根本起不来,报 NoResourceAvailableException。

一句话记住:TaskManager 是一台干活的机器(JVM 进程),Slot 是这台机器上切出来的一个带独立内存配额的工位。 后面要反复排查的「某些 TM 被打满、某些几乎闲置」,根子多半就在 Slot 的划分与摆放没配好。

4. 从代码到执行:三次图变换

角色就位,还差一条链路:你写的 DataStream 代码,到集群里真正跑起来的物理任务之间,隔着三次图变换。看懂这三张图,才算真正理解「我这几行代码最后变成了几个线程、分别跑在哪」。

Flink 从代码到物理执行的三次图变换图。第一层 StreamGraph(逻辑 DAG):Source、map、keyBy/window、Sink 四个逻辑算子按顺序连成有向图,是 DataStream API 直接翻译出的原始结构。第二层 JobGraph(算子链合并后,由 Client 提交给 JobManager):把 forward 且并行度相同的相邻算子 Source 和 map 合并成一个算子链节点 [Source → map],减少节点数;window 和 Sink 作为独立 JobVertex。第三层 ExecutionGraph(按并行度展开的物理执行图):每个 JobVertex 按并行度展开成多个并行 subtask,图中 [Source→map] 和 window 各展开为 2 个 subtask、Sink 为 1 个 subtask,[Source→map] 到 window 之间因 keyBy 是 hash 全连接重分区(all-to-all),window 到 Sink 之间为收敛连接。层与层之间用向下箭头标注 chaining 合并算子链、按并行度展开 subtask。

三次图变换: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:

三张图一句话概括:StreamGraph 是你写了什么,JobGraph 是合并优化后提交了什么,ExecutionGraph 是集群里实际并行跑成了什么。 并行度正是从 JobGraph 放大到 ExecutionGraph 这一步的关键旋钮,它直接决定你需要多少 Slot、每个算子摊到几个线程。

5. 算子链与 Slot 共享

上一节里一带而过的 chaining,值得单独展开——它和 Slot 共享一起,决定了同样一份 ExecutionGraph 最终会以怎样的密度落在集群上。

5.1 Operator Chaining:为什么要合并算子

算子链是 Flink 默认开启的核心性能优化:把满足条件的相邻算子塞进同一个线程里串行执行,一条数据在链内算子之间直接以方法调用传递,而不必走「序列化 → 网络或内存队列 → 反序列化」这一圈。

合并带来三重收益:

合并的条件(简化说)是:相邻算子并行度相同、连接方式为 forward(一对一)、处在同一个 slot sharing group,且 chaining 没有被显式关掉。上图里的 Source → map 就是典型——并行度同为 2、forward 直连,于是合并成一条链;而 map → window 之间夹着 keyBy(hash 重分区而非 forward),链在这里断开。

5.2 什么时候要主动断链

默认合并总体是好事,但有两种场景值得主动断链:

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 里)。

这带来两个好处:

  1. Slot 总数只需等于最大算子并行度,而不是所有算子并行度之和。假设作业里 source 并行度 2、window 并行度 8、sink 并行度 2,没有 slot sharing 就要 2+8+2=12 个 Slot;有了 slot sharing,一个 Slot 里可以同时放下 source、window、sink 各一个 subtask,只需 8 个 Slot(即最大并行度)。
  2. 轻重算子自动搭配:一个 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」,它直接影响数据分布、状态正确性以及会不会倾斜。最常用的是下面四种。

Flink subtask 间四种数据分区策略示意图。四个小面板各画上游 subtask 到下游 subtask 的数据流向。左上 forward(一对一·不重分区):上游 1→下游 1、上游 2→下游 2,一一直连,要求上下游并行度相同,是算子链合并的前提。右上 hash / keyBy(按 key 哈希重分区):上游两个 subtask 到下游三个 subtask 全连接,含义是按 key 的哈希值路由,保证相同 key 的记录总落到同一个下游 subtask。左下 rebalance(轮询均衡·解决数据倾斜):上游两个 subtask 以轮询方式把记录均匀分发到下游三个 subtask,用于打散数据倾斜。右下 broadcast(广播·每条复制到全部):上游每个 subtask 的每条记录都复制发送到所有下游 subtask,用于广播小维表或规则。

四种分区策略:forward 一对一(算子链前提)、hash/keyBy 按 key 哈希(同 key 同 subtask,keyed state 的基础)、rebalance 轮询均衡(治倾斜)、broadcast 广播(每条复制到全部,用于广播维表/规则)。

注:回到广告链路,曝光计数必须用 keyBy(campaignId)(hash,同 campaign 同 subtask 才能正确累加);若少数超大 campaign 造成倾斜、而某一步又不需要按 key,就用 rebalance 把负载摊平;广告主预算阈值这类小表则用 broadcast 发给所有 subtask,供本地 pacing 判断。选错分区策略的后果只有两种:该 keyBy 却 rebalance 会算错,该 rebalance 却放任热 key 会压垮单个 subtask。

7. 部署模式:Session、Per-Job 与 Application Mode

同一个作业、同一张 ExecutionGraph,还可以用不同的部署模式跑起来,区别在于集群的生命周期与**main() 在哪里执行**,而这两点又直接决定了资源隔离粒度和日常运维方式。

底层资源框架上,以上模式都可以跑在不同的资源管理器之上:

生产建议:核心实时作业用 Application Mode + on K8s——强隔离(一个作业崩了不影响别人)、main() 在集群内执行、还能弹性伸缩;Session 模式留给那些数量多、体量小、能容忍互相影响的作业。

8. 一个作业从提交到运行的完整流程

最后把前七节的角色、图变换、分区策略与部署模式串成一条线,完整走一遍从 env.execute() 到数据开始流动的生命周期:

  1. Client 编译:执行 main(),先生成 StreamGraph,再做算子链合并得到 JobGraph(Application 模式下这一步在集群内完成)。
  2. 提交作业:Client 把 JobGraph 连同用户 jar 与依赖一起提交给 Dispatcher。
  3. 拉起 JobMaster:Dispatcher 为这个作业启动一个专属 JobMaster,并把 JobGraph 交给它。
  4. 展开 ExecutionGraph:JobMaster 按各算子并行度把 JobGraph 展开成 ExecutionGraph(一堆 subtask),同时算出需要多少 Slot。
  5. 申请资源:JobMaster 向 ResourceManager 申请 Slot;跑在 YARN 或 K8s 上时,Slot 不够则由 ResourceManager 动态拉起新的 TaskManager。
  6. 分配 Slot:ResourceManager 在考虑 slot sharing 的前提下,把空闲 Slot 分配给 JobMaster。
  7. 部署 subtask:JobMaster 把各个 subtask 部署到对应 TaskManager 的 Slot 上,TM 随即启动 task 线程。
  8. 建立数据通道并运行:subtask 之间按分区策略(forward / hash / rebalance / broadcast)建立数据传输通道,source 开始从 Kafka 拉数据,整条流水线持续处理。
  9. 常驻运行与容错:运行期间 JobMaster 周期性触发 checkpoint 并监控 task 心跳,task 失败则触发 restart、从最近一次 checkpoint 恢复状态(详见 状态与 Checkpoint)。

注:这条链路上任何一环卡住,都会如实反映到 Web UI 上——卡在第 6 步报 NoResourceAvailableException 就是 Slot 不够;跑起来之后部分 subtask 忙、部分闲,则是并行度、分区策略或 slot sharing 没配好。地基到这里就铺完了,下一篇会回到 §2.3 那个被一笔带过的 window,讲清楚事件时间究竟怎么推进、迟到数据还能等多久。

参考


–views
Share this post on:

Previous Post
Flink 深挖 · 时间语义、窗口与 Watermark
Next Post
Apache Flume:从 Source·Channel·Sink 到广告日志不丢链路与选型