这是什么
这是一份为有 Flink 使用经验的开发者准备的源码深入学习课程。
市面上多数教程停留在 API 层面——教你 flatMap、keyBy、window 怎么用。但当你把 Flink 用到生产环境,一定会遇到 API 文档回答不了的问题:
"为什么我的作业启动这么慢?" "为什么 keyBy 之后算子不能被 Chain 到一起?" "反压到底是怎么传播的?怎么排查?" "Checkpoint 卡住了,是哪个环节的问题?" "K8s 上 TM Pod 一直没有被创建,卡在哪一步了?"
这些问题只能靠读源码来回答。本课程的目标就是让你具备独立阅读 Flink 源码并定位问题的能力。
通过 10 个阶段、49 节课、每节课 10-20 分钟的节奏,用层层递进的方式,带你从 StreamExecutionEnvironment 的构造开始,一路深入到 Task 线程的执行循环、Checkpoint 的分布式快照、K8s/YARN 的容器部署。
课程适合谁
| 写过 DataStream 或 Table/SQL 作业 | |
本课程假设你已经:
用 DataStream API 或 SQL 写过 Flink 作业 知道 JobManager / TaskManager 是什么(概念层面即可) 跑过 flink run,配过flink-conf.yaml
本课程不假设:
你读过任何 Flink 源码 你理解 ExecutionGraph / JobGraph / StreamGraph 的区别(这正是课程要讲的)
如果你完全没有接触过 Flink,建议先去官方文档跑一遍 DataStream API 教程 或 Table API 教程,再回来学习本课程。大数据组件的源码学习门槛是客观存在的——先有使用体感,再看源码才有共鸣。
课程解决的核心问题
市面上很多 Flink 教程教你怎么用 API(flatMap、keyBy、window),但当你遇到:
"为什么我的作业启动这么慢?" "为什么 keyBy 之后算子不能被 Chain?" "反压到底是怎么传播的?" "Checkpoint 卡住了怎么办?" "K8s 上 TM Pod 为什么没有被创建?"
这些问题只看 API 文档是解决不了的,需要理解内部机制。本课程就是为回答这些问题而设计的。
设计理念
1. 剥洋葱式学习
用户代码(你写的 WordCount) ↓ API 层做了什么(惰性求值、Transformation DAG) ↓ 图转换做了什么(StreamGraph → JobGraph → ExecutionGraph) ↓ 调度与部署做了什么(JM → Scheduler → Slot → TM) ↓ 运行时执行了什么(Task → StreamTask → Operator) ↓ 数据怎么流转的(RecordWriter → ResultPartition → InputGate) ↓ 状态怎么管理的(Checkpoint → StateBackend → KeyGroup) ↓ 集群怎么部署的(K8s Pod / YARN Container → ActiveResourceManager)每一层只关心一层的事,不跳层、不混杂。学完一层再剥开下一层。
2. 一个例子贯穿始终
全部 49 节课都用 WordCount(Source → flatMap → keyBy → sum → print)作为线索。你不用每次切换上下文,始终在同一份代码中越挖越深。
3. 先概念、后代码、再图示
每节课的结构:
是什么:这个概念解决什么问题(2-3 分钟) 怎么实现:带着你看关键源码,逐行解释(10-12 分钟) 图总结:一张图把知识点串起来(2-3 分钟)
4. 聚焦核心路径,不追求全面
Flink 有 200+ 万行 Java 代码,你不可能读完 本课程只讲核心链路(约占代码量的 15-20%),但覆盖了 80% 的实际问题场景 掌握核心链路后,你自己就能读其他模块的代码
学习方法论
建议节奏
每天 2-3 节课 = 每天 30-45 分钟核心路径(阶段一→六,27 节课):约 2 周部署路径(阶段十,6 节课):约 3 天进阶路径(阶段七→八,10 节课):约 1 周可选路径(阶段九,6 节课):按兴趣怎么读最有效
打开 IDE:在 IntelliJ IDEA 中打开 Flink 工程,跟着课程中的类名和行号跳转 打断点:在关键的 transform()、isChainable()、processElement()方法上打断点,运行 WordCount 看调用栈每阶段完成时回顾:每个阶段最后一课是回顾总结,把零散知识点串起来 不求一次全懂:第一遍不求完全理解,先建立"是什么"的认知;第二遍再看"为什么这样设计"
10 个阶段的逻辑关系
┌───────────────────────────────────────────┐ │ Flink 源码学习全景图 │ └───────────────────────────────────────────┘ 前置课:建立心智模型(Flink 是什么、JM/TM/Slot/Barrier 基本概念) │ 阶段一:StreamExecutionEnvironment —— 大门(环境发现、惰性求值) │ 阶段二:Transformation DAG —— 用户 API 的背后(flatMap/keyBy/sum) │ 阶段三:StreamGraphGenerator —— 第一个图转换(节点+边) │ 阶段四:StreamingJobGraphGenerator —— 算子链优化(Chain 算法) │ 阶段五:JobMaster & ExecutionGraph —— 并行化展开与调度 │ ├──────────── 核心路径完成 ────────────┐ │ │ ▼ ▼ 阶段六:Task & StreamTask 阶段十:K8s/YARN 部署 (执行循环、Mailbox 模型) (Pod/Container 生命周期) │ │ ▼ │ 阶段七:网络栈 ◀──────────┘ (数据流转、反压) │ ▼ 阶段八:Checkpoint & State (容错恢复、状态管理) │ ▼ 阶段九:Table/SQL 引擎 (Calcite 集成、优化规则)课程目录一览
getExecutionEnvironment()、构造函数、execute() | ||
flatMap、keyBy、sum 的幕后 | ||
generate()、transform()、translateInternal() | ||
isChainable()、ChainingStrategy、createChain() | ||
TaskExecutor、StreamTask、AbstractStreamOperator | ||
RecordWriter、ResultPartition、反压机制 | ||
StateBackend、增量快照 | ||
ActiveResourceManager |
与官方文档的关系
官方文档教你怎么用 Flink(API 参考、配置说明) 本课程教你Flink 是怎么做的(内部机制、源码实现) 两者互补:先用官方文档写出能跑的代码,再用本课程理解背后的原理
源码版本
本课程基于 Apache Flink 源码仓库,代码引用指向 src/main/java/org/apache/flink/ 下的文件。Flink 的类名和方法名在版本间保持稳定,核心逻辑变化很小,因此适用于 Flink 1.15+ 版本。
开始之前
用 IntelliJ IDEA 打开 Flink 工程( pom.xml在根目录)在 flink-examples/flink-examples-streaming/中找到WordCount.java运行一次,确认环境正常
夜雨聆风