跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
🌊 Flink 教程 · 2 / 7 查看系列简介 →

Flink 教程(二):架构与原理

文章目录
  1. 2.1 Flink 架构
  2. 2.1.1 集群组件
  3. 2.1.2 Task Slot 和资源
  4. 2.1.3 核心架构详解
  5. 2.1.4 集群模式
  6. 2.2 有状态流处理
  7. 2.2.1 什么是状态
  8. 2.2.2 Checkpoint 机制
  9. 2.2.3 非对齐 Checkpoint
  10. 2.2.4 内存管理与反压
  11. 2.3 时间语义与水位线
  12. 2.3.1 时间概念
  13. 2.3.2 水位线(Watermark)
  14. 2.3.3 延迟处理
  15. 2.4 状态后端(State Backends)
  16. 2.4.1 状态后端对比详解
  17. 2.5 Savepoint
  18. 2.6 窗口机制
  19. 2.6.1 窗口触发时间
  20. 2.6.2 窗口类型
  21. 2.6.3 算子与窗口的关系
  22. 2.7 动态表

📌 本文是「Flink 教程」系列第 二 篇 · 系列目录

Flink 是典型的 Master-Slave (主从) 架构,类似于 ” 包工头与工人 ” 的关系。

2.1.1 集群组件

Flink 运行时由两种类型的进程组成:JobManager 和 TaskManager。

JobManager 组件:

组件职责
ResourceManager资源提供、回收、分配,管理 task slots
Dispatcher提供 REST 接口提交作业,运行 WebUI
JobMaster负责管理单个 JobGraph 的执行

TaskManager:

  • 执行作业流的 task
  • 缓存和交换数据流
  • 资源调度的最小单位是 task slot
  • 每个 slot 代表 TaskManager 资源的固定子集

2.1.2 Task Slot 和资源

Slot 数量:

  • 每个 TaskManager 至少有一个 slot
  • Slot 数量表示并发处理 task 的数量
  • 一个 slot 中可以执行多个算子(通过算子链)

Slot 共享:

特性说明
共享规则同一作业的不同 Task 的 SubTask 可以共享 Slot
最大共享整个作业管道可以放在一个 Slot 中执行
核心优势资源利用率更高,无需计算总 Task 数

Slot Sharing 示例 Source(1) → Map(1) → Sink(1),这 3 个 SubTask 可共享 1 个 Slot

2.1.3 核心架构详解

组件角色

组件角色职责
JobManagerMaster/包工头作业调度、Checkpoint 协调、故障恢复、管理 Resource Manager 和 Dispatcher
TaskManagerWorker/工人实际执行计算的 JVM 进程,每个 TM 包含一定数量的 Slot
Slot工位TM 资源的固定子集(主要是内存隔离,CPU 共享),决定并行执行数
SubTask子任务算子根据并行度复制出来的物理执行单元(线程)

核心机制

  • Slot Sharing (槽位共享): 允许同一个 Job 的不同 Task 的 SubTask 共享同一个 Slot(如 Source + Map + Sink 可在一个 Slot 跑),提高资源利用率。
  • 通信:
    • 控制流: Akka (RPC) —— JM 与 TM 之间指令交互
    • 数据流: Netty —— 跨 TM 的数据传输(Shuffle)
graph TD
    subgraph "客户端"
        Client
    end

    subgraph "JobManager"
        Dispatcher[Dispatcher<br/>WebUI]
        JobMaster[JobMaster<br/>作业管理]
        RM[ResourceManager<br/>资源管理]
    end

    subgraph "TaskManager 集群"
        TM1["TaskManager 1<br/>Slot: ★ ★ ★"]
        TM2["TaskManager 2<br/>Slot: ★ ★ ★"]
        TMn["TaskManager N<br/>Slot: ★ ★ ★"]
    end

    Client --> Dispatcher
    Dispatcher --> JobMaster
    JobMaster --> RM
    RM <-->|"Akka RPC<br/>控制流"| TM1
    RM <-->|"Akka RPC<br/>控制流"| TM2
    RM <-->|"Akka RPC<br/>控制流"| TMn

    TM1 <-->|"Netty<br/>数据交换"| TM2
    TM2 <-->|"Netty<br/>数据交换"| TMn
    TM1 <-->|"Netty<br/>数据交换"| TMn

    subgraph "TM1 内部"
        S1[Slot 1: Source→Map→Sink]
        S2[Slot 2: Filter]
        S3[Slot 3: Agg]
    end

2.1.4 集群模式

模式生命周期资源隔离适用场景
Session 集群长期运行,作业共享资源竞争交互式分析、短作业
Job 集群作业启动时创建作业隔离大型作业、长时运行
Application 集群随应用启动最佳隔离云原生部署

2.2 有状态流处理

有状态流处理是 Flink 的核心概念,指在处理事件流时维护和更新状态信息。

2.2.1 什么是状态

在数据流中,许多操作(如窗口操作)需要记住跨多个事件的信息,这些操作称为有状态操作。

状态类型:

// Keyed State - 按 Key 分区的状态
ValueState<T>              // 单值状态,存储单个值
ListState<T>               // 列表状态,存储元素列表
MapState<K, V>             // 映射状态,存储键值对
ReducingState<T>           // 聚合状态,自动聚合
AggregatingState<IN, OUT>  // 聚合状态,自定义聚合

// Operator State - Operator 级别的状态(无 Key 分组)

Keyed State 的特点:

  • 仅在 keyed streams 上可用(即 keyBy 操作之后)
  • 状态与特定的 key 绑定
  • 访问仅限于当前事件的 key
  • 状态更新是本地操作,保证一致性

2.2.2 Checkpoint 机制

Flink 使用检查点和流重播的组合来实现容错。

Checkpoint 原理:

  1. 在每个输入流的特定点标记检查点
  2. 同时保存每个操作符的对应状态
  3. 从检查点恢复时,重置状态并从该点重放记录

Barrier(屏障):

  • 屏障注入数据流,作为数据的一部分流动
  • 屏障不会超过记录,它们严格按顺序流动
  • 一个屏障分隔当前快照和下一个快照的记录

Checkpoint 流程:

  1. 屏障从源头注入(记录在 Kafka 中的偏移量)
  2. 屏障向下游流动
  3. 中间操作符从所有输入流收到屏障后,发送屏障到输出流
  4. Sink 操作符收到所有输入的屏障后,确认检查点完成

Checkpoint 屏障流动图

sequenceDiagram
    participant Source
    participant Map
    participant KeyBy
    participant Sink

    Note over Source: Barrier n 注入<br/>记录 Kafka 偏移量
    Source->>Map: 数据 + Barrier n
    Map->>KeyBy: 数据 + Barrier n
    KeyBy->>Sink: 数据 + Barrier n

    Note over Sink: 收到所有输入 Barrier n<br/>确认 Checkpoint 完成 ✓

Checkpoint vs Savepoint 对比

特性CheckpointSavepoint
触发自动(周期性)手动
生命周期任务停止即删(默认)永久保留
用途活下来(防挂)活得更好(升级)

2.2.3 非对齐 Checkpoint

Flink 1.11 引入了非对齐 Checkpoint,基本思想是检查点可以超过所有飞行中的数据,只要飞行中的数据成为操作符状态的一部分。

特点:

  • 确保屏障以最快速度到达 sink
  • 特别适合慢速数据路径的场景
  • 避免了长时间的对齐等待

非对齐 Checkpoint:解决严重反压导致 Checkpoint 超时的问题,Barrier 插队优先处理,但需将缓冲区数据一起存盘

2.2.4 内存管理与反压

TaskManager 内存估算

方式推荐风险
增量聚合使用 Reduce/Aggregate Function内存只存累加结果,占用极小
全量缓存使用 ProcessWindowFunction内存存原始数据 List,极大可能 OOM

估算公式:Key 数量 × 状态对象大小(需考虑 Java 对象头开销)

反压 vs Watermark

对比项反压Watermark 不一致
本质物理现象,下游处理慢逻辑现象,水位线滞后
表现缓冲区填满,堵塞上游窗口无法关闭,状态堆积
传播方向下游 → 上游(反向)上游 → 下游(正向)

关系图示:

flowchart TD
    subgraph "正常情况"
        A1[Source] -->|数据| B1[Map]
        B1 -->|数据| C1[KeyBy]
        C1 -->|数据| D1[Sink]
        D1 -.->|"✓ 处理正常"| B1
    end

    subgraph "反压场景"
        A2[Source] -->|数据| B2[Map]
        B2 -->|数据| C2[KeyBy]
        C2 -->|"⚠️ 缓冲区满"| D2[Sink]
        D2 -.->|"🚫 反压 Backpressure"| C2
        C2 -.->|"🚫 反压"| B2
        B2 -.->|"🚫 反压"| A2
    end

    style D2 fill:#ffcccc
    style C2 fill:#ffdddd
    style B2 fill:#ffeeee

核心关系:

  • 反压 → Watermark:反压会导致缓冲区积压,Watermark 传递变慢
  • Watermark → 反压:Watermark 滞后导致状态无限堆积,最终引发反压

2.3 时间语义与水位线

2.3.1 时间概念

Flink 支持三种时间概念:

时间类型说明特点
事件时间(Event Time)事件实际发生的时间结果可重现,与处理时间无关
处理时间(Processing Time)系统处理事件的时间简单,但结果不可重现
摄入时间(Ingestion Time)Flink 读取事件的时间介于两者之间

处理时间 vs 事件时间:

处理时间是最简单的时间概念,不需要协调流和机器之间的协调,提供最佳性能和最低延迟。但在分布式和异步环境中,处理时间不提供确定性,因为它受记录到达系统的速度影响。

事件时间处理会产生完全一致和确定性的结果,无论事件何时到达或如何排序。然而,除非已知事件按顺序到达,否则事件时间处理会在等待乱序事件时产生一些延迟。

2.3.2 水位线(Watermark)

水位线是 Flink 中测量事件时间进度的机制。

水位线的作用:

  • 声明事件时间到达某个时间点
  • 标记该时间之前的事件应该都已到达
  • 触发窗口关闭和计算

水位线策略:

WatermarkStrategy<Event> strategy = WatermarkStrategy
    .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(20))  // 最大无序边界
    .withTimestampAssigner((event, timestamp) -> event.timestamp);

DataStream<Event> withWatermarks = stream.assignTimestampsAndWatermarks(strategy);

Watermark(t) 声明:

  • 声明事件时间已到达 t
  • 不应有更多时间戳 t' <= t 的元素
  • 即时间戳小于等于水印的事件都应该已到达

乱序流中的水位线:

  • 一般来说,水位线是对某个时间点之前所有事件都已到达的声明
  • 一旦水位线到达操作符,操作符可以将其内部事件时间时钟推进到水位线的值

Watermark 多并行度细节

  • Watermark 在 Source Operator 中生成,并且在每个 Source Operator 的子 Task 中都会独立生成 Watermark
  • Watermark 会逐步更新下游算子中的 Watermark 水位线
  • 如果多个 Watermark 同时更新一个算子 Task 的当前事件时间,Flink 会选择最小的水位线来更新
  • 当一个 Window 算子 Task 中水位线大于了 Window 结束时间,就会立即触发窗口计算
  • 慢分区会拖慢整体进度:如果某个分区 Watermark 落后太多,可能是该分区数据积压或处理慢,建议监控各并行度的 Watermark 进度

水位线核心公式

  • 公式:Watermark = MaxEventTime - MaxOutOfOrderness(最大允许乱序时间)
  • 传播:多对一场景下遵循木桶效应,下游取所有上游输入中最小的 Watermark

木桶效应图示:

flowchart LR
    subgraph "上游 Source (并行度=3)"
        P1["Partition 1<br/>WM: 12:05"]
        P2["Partition 2<br/>WM: 11:58"]
        P3["Partition 3<br/>WM: 12:02"]
    end

    subgraph "下游 Window"
        W["Window<br/>取最小值: 11:58"]
    end

    P1 -->|"WM: 12:05"| W
    P2 -->|"WM: 11:58<br/>⚠️ 最慢分区"| W
    P3 -->|"WM: 12:02"| W

    style P2 fill:#ffcccc,stroke:#ff0000
    style W fill:#ffdddd,stroke:#ff0000

公式详解:

Watermark = 当前分区最大事件时间 - 最大允许乱序时间

例如:最大事件时间 12:05,最大无序边界 5 分钟
     Watermark = 12:05 - 5min = 12:00
     → 声明:12:00 之前的数据都已到达

2.3.3 延迟处理

允许的延迟(Allowed Lateness):

  • 延迟事件是那些在水位线之后到达的事件
  • 可以配置允许延迟的时间窗口
  • 在延迟窗口内,延迟事件仍会被处理

侧输出(Side Output):

  • 将被丢弃的延迟事件发送到侧输出流
  • 用于后续单独处理延迟数据

迟到数据处理(三道防线)

  1. Watermark 自带缓冲:设置 MaxOutOfOrderness(如 5s)
  2. AllowedLateness(窗口宽容期):窗口触发后保留状态一段时间(如 1min),迟到数据会导致窗口再次触发
  3. Side Output(侧输出流):彻底迟到的数据兜底存入侧输出流,人工处理

2.4 状态后端(State Backends)

常见误区 state.backend: filesystem 是 Flink 1.13 之前的旧写法,等同于 HashMapStateBackend,运行时数据在 JVM 堆内存,Checkpoint 时才写入文件系统。

Flink 提供三种状态后端:

状态后端存储位置适用场景
MemoryStateBackendJVM 堆内存开发测试、小状态
FsStateBackend文件系统 + 堆内存中等状态、生产环境
RocksDBStateBackendRocksDB + 文件系统大状态、生产环境

2.4.1 状态后端对比详解

类型场景存储位置配置
HashMapStateBackend低延迟、状态量可控(< 几 GB)TM 堆内存state.backend: hashmap
EmbeddedRocksDBStateBackend超大状态(TB 级)、长窗口本地磁盘(RocksDB)+ 内存缓存state.backend: rocksdb

💡 Tips: RocksDB JNI 传输限制(生产环境重要)

  • RocksDB 通过 JNI 的方式进行数据交互
  • 每次能够传输的最大数据量为 2^31 字节(约 2GB)
  • 在 RocksDBStateBackend 合并的状态数据量大小不能超过此限制
  • 否则将会导致状态数据无法同步,这是 RocksDB 采用 JNI 方式的限制
  • 用户在使用过程中应当注意,避免单条状态记录超过 2GB

状态后端选择建议:

场景推荐后端原因
开发调试MemoryStateBackend快速、简单
小状态(<1GB)FsStateBackend性能好
大状态(>1GB)RocksDBStateBackend容量大
状态 > 内存RocksDBStateBackend磁盘存储

2.5 Savepoint

Savepoint 是手动触发的检查点,用于:

  • 作业停止和恢复
  • 应用程序升级
  • 集群迁移
  • A/B 测试

特点:

  • 手动触发,不会自动过期
  • 依赖常规检查点机制
  • 从 savepoint 恢复时进行精确的状态恢复

Savepoint 要点:

  • 必须为有状态算子设置 UID(.uid("name")),否则代码修改后状态无法恢复

2.6 窗口机制

2.6.1 窗口触发时间

窗口触发条件:Watermark ≥ WindowEndTime

窗口触发时间线:

flowchart LR
    subgraph "数据到达阶段"
        A1((12:00<br/>窗口开启))
        A2(12:02<br/>事件A)
        A3(12:05<br/>事件B)
        A4(12:08<br/>事件C)
    end

    subgraph "等待水位线阶段"
        B1((12:10<br/>窗口关闭时间点))
        B2{等待<br/>Watermark}
    end

    subgraph "窗口触发阶段"
        C1((12:15<br/>Watermark>=12:10<br/>窗口计算触发))
    end

    A1 --> A2 --> A3 --> A4 --> B1 --> B2 --> C1

    style A1 fill:#e1f5fe
    style C1 fill:#c8e6c9
    style B2 fill:#fff3e0

关键公式:

看到结果时刻 = 窗口结束时间 + Watermark(最大无序边界)

示例:窗口 [12:00, 12:10),MaxOutOfOrderness=5min
     必须等到 Watermark >= 12:10
     即:12:05 之后的数据最大时间戳 >= 12:05
     最早 12:10 看到结果,最迟 12:15 看到结果

2.6.2 窗口类型

类型特点风险
滚动窗口(Tumbling)数据只属于 1 个窗口-
滑动窗口(Sliding)数据会属于 N 个 窗口资源消耗大(数据逻辑膨胀)

2.6.3 算子与窗口的关系

  • 并非 1 个窗口占用 1 个算子
  • 1 个 WindowOperator (SubTask) 管理着该 Slot 内所有 Key 的成千上万个窗口对象
  • Slot 之间隔离,但 Slot 内部的窗口共享内存(大户 Key 会导致 OOM)

2.7 动态表

Flink 的 Table API 和 SQL 基于动态表概念工作。

核心原理:

  • 流可以转换为表(Changelog 流上的查询)
  • 表可以转换回流(输出)
  • 支持连续查询(Continuous Query)

流与表的转换:

流 --(asTable)--> 动态表 --(toChangelogStream)--> 流

Append-Only Stream  <-->  动态表  <-->  Retract Stream
Upsert Stream        <-->  动态表  <-->  Upsert Stream

参见: 01-应用场景 | 03-部署与运维 | 04-Java开发实践 | 附录


标签: #flink #架构 #原理 #大数据 #流处理


RAG 智能问答

针对本文继续提问:《Flink 教程(二):架构与原理》