乐于分享
好东西不私藏

Flink源码学习系列课程:前置课,Flink核心概念

Flink源码学习系列课程:前置课,Flink核心概念

说明

本课程假设你已经用过 Flink。这一课不是从零教你 Flink,而是在进入源码之前,统一概念说法和命名,确保后续课程中每个术语的指代都是明确的。

如果你对以下概念已经很熟悉,直接跳过即可。

Flink 是什么

Apache Flink 是一个分布式有状态流处理引擎

拆解这三个关键词:

1. 流处理(Stream Processing)

数据来一条处理一条,而不是等一批数据凑齐了一起处理。

批处理(Batch):等所有数据到了再处理  [data1, data2, data3, data4, data5] → 一次性处理 → [result1, result2, ...]流处理(Stream):来一条处理一条  data1 → 处理 → result1  data2 → 处理 → result2  data3 → 处理 → result3  ...实时产生结果

为什么重要:实时风控、实时大屏、实时推荐等场景,不能等数据攒够一批再处理。

2. 分布式(Distributed)

一个 Flink 作业可以运行在多台机器上,每台机器处理一部分数据。

                    ┌──────────────┐                    │  JobManager  │  ← 老板:协调调度                    │   (1 个)     │                    └──────┬───────┘               ┌───────────┼───────────┐               ▼           ▼           ▼        ┌──────────┐ ┌──────────┐ ┌──────────┐        │TaskManager│ │TaskManager│ │TaskManager│  ← 工人:执行任务        │  (多个)   │ │  (多个)   │ │  (多个)   │        └──────────┘ └──────────┘ └──────────┘
  • JobManager (JM) :一个作业一个 JM,负责协调、调度、Checkpoint 管理
  • TaskManager (TM) :执行具体任务,通常有很多个,每个 TM 有多个 Slot(执行槽位)

3. 有状态(Stateful)

Flink 能记住之前处理过的数据

// 例如 WordCount 中的 sum 操作:// 看到 "hello" 第 1 次 → count = 1// 看到 "hello" 第 2 次 → count = 2(需要记住之前的 1)// 看到 "hello" 第 3 次 → count = 3(需要记住之前的 2)// 这个"记住的值"就是状态(State)

Flink 的状态是容错的:即使机器挂了,状态也能恢复,不会丢数据。

一个 Flink 作业的基本结构

Source(数据源) → Transform(处理) → Sink(输出)    ↓                  ↓                ↓  读数据            转换/聚合         写出去 (Kafka/文件)    (flatMap/sum)    (数据库/文件)

以 WordCount 为例:

文件 → 拆词("hello world" → ["hello","world"]) → 按词分组 → 计数(hello:2) → 打印

Flink 的执行流程(宏观视角)

1. 用户写代码(DataStream API)      ↓2. Flink 构建执行计划(Transformation → StreamGraph → JobGraph → ExecutionGraph)      ↓3. JobManager 调度,把任务分发给 TaskManager      ↓4. TaskManager 执行任务(开线程,跑算子)      ↓5. 数据在算子之间流动,产生结果

本教程就是带你从第 1 步到第 4 步,把每一个环节的源码看清楚。

关键概念速查表

概念
一句话解释
StreamGraph
API 调用产生的逻辑执行图
JobGraph
优化后的执行图(做了算子链合并)
ExecutionGraph
并行化后的调度图(加入了运行时状态)
Task / StreamTask
TM 上执行的最小单元
Operator / StreamOperator
具体处理数据的算子(flatMap, sum...)
Slot
TM 上的资源单位,一个 Slot 跑一个 Task
Chain
把多个算子合并到一个 Task 中执行(优化)
Checkpoint
分布式快照,用于故障恢复
StateBackend
状态的存储方式(内存 or RocksDB)
Barrier
数据流中的"分隔标记",触发 Checkpoint

接下来

下一课我们从 StreamExecutionEnvironment.getExecutionEnvironment() 开始,看 Flink 如何创建一个执行环境。