Skip to content
Charles Shao
Go back

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

views

广告实时链路里,Kafka 只是把事件搬到了流上(见 Kafka 核心原理);真正把这些曝光、点击、转化事件有状态地、低延迟地、精确一次地算成钱和报表的,是流处理引擎。这个系列专门深挖 Apache Flink——从运行时架构、时间与窗口、端到端 Exactly-Once、DataStream,到状态容错、反压调优与生产部署,一路钻到能上生产、能定位问题的深度。开篇这一篇先把地基铺平:Flink 的流处理模型长什么样,一个作业提交下去,在集群里到底是怎么跑起来的。

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

TL;DR

Table of contents

Open Table of contents

广告系统的实时链路——曝光/点击/转化事件从 Kafka 流进来,要实时算出:每个广告主花了多少钱(计费)、预算还剩多少(pacing 控速)、点击归因到哪次曝光(attribution)、每个创意的实时 CTR(报表)——对计算引擎的要求非常苛刻:延迟要低(pacing 慢一步就超投)、结果要准(计费不能算错一分钱)、还要能扛住乱序和迟到的事件。为什么最后大家选了 Flink,而不是 Spark Streaming 或 Kafka Streams?看四个维度。

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

Spark Streaming(含 Structured Streaming 的微批模式)本质是”攒一批再算”:它把流切成一个个小批次(micro-batch,比如每 1 秒一批),每批当成一个小的批作业跑。这带来两个固有代价:

Flink 是逐条处理的真流:一条曝光事件进来,立刻流过 source → map → keyBy → window → sink 这条算子链,中间不攒批。延迟是毫秒级,且延迟不随吞吐线性恶化。对 pacing 这种”预算快花完了要立刻降速”的场景,毫秒 vs 秒级是决定性的。

注: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 才到达。如果按”到达时间(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 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 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 客户端负责构建 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 就各拿 1/3,一个 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,减少节点数、省掉算子间的序列化和线程切换。

上图里 Sourcemap 被合并成一个链节点 [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 共享

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 数只需等于最大并行度

Flink 默认还有一个巧妙机制:Slot 共享(Slot Sharing)——同一个作业里、不同算子的 subtask,可以共享同一个 Slot(默认都在同一个 slot sharing group default)。

它带来两个好处:

  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,等于给重算子”包间”。这既是隔离资源的手段,也是修复”重算子挤一起”的关键动作之一。

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 广播(每条复制到全部,用于广播维表/规则)。

AdTech 例子:曝光计数用 keyBy(campaignId)(hash,保证同 campaign 同 subtask 才能正确累加);若发现少数超大 campaign 造成倾斜、且某一步不需要按 key,就用 rebalance 摊平;把”广告主预算阈值配置”这类小表用 broadcast 发给所有 subtask 做本地 pacing 判断。选错分区策略要么算错(该 keyBy 却 rebalance),要么倾斜(该 rebalance 却让热 key 压垮一个 subtask)。

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

同样一个作业,可以用不同的部署模式跑,区别在于集群的生命周期和**main() 在哪执行**。这决定了资源隔离粒度和运维方式。

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

生产建议:核心实时作业用 Application Mode + on K8s——强隔离(一个作业崩不影响别人)、main() 在集群跑、弹性伸缩。Session 模式留给”一堆小作业、能容忍互相影响”的场景。

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

把前面的角色和图变换串起来,走一遍env.execute() 到数据开始流动的完整生命周期:

  1. Client 编译:执行 main(),把 DataStream 描述先生成 StreamGraph,再做算子链合并生成 JobGraph(Application 模式下这步在集群内做)。
  2. 提交作业:Client 把 JobGraph(+ 用户 jar、依赖)提交给 Dispatcher
  3. 拉起 JobMaster:Dispatcher 为这个作业启动一个专属 JobMaster,把 JobGraph 交给它。
  4. 展开 ExecutionGraph:JobMaster 把 JobGraph 按各算子并行度展开成 ExecutionGraph(一堆 subtask),并计算需要多少 Slot。
  5. 申请资源:JobMaster 向 ResourceManager 申请 Slot;on YARN/K8s 时,ResourceManager 若 Slot 不够会动态拉起新的 TaskManager
  6. 分配 Slot:ResourceManager 把空闲 Slot(考虑 slot sharing)分配给 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 配置问题。

参考


views
Share this post on:

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