说明
本课程假设你已经用过 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 步,把每一个环节的源码看清楚。
关键概念速查表
接下来
下一课我们从 StreamExecutionEnvironment.getExecutionEnvironment() 开始,看 Flink 如何创建一个执行环境。
夜雨聆风