Flink internals ←STREAM PROCESSING / SOURCE READING / LESSON 00
Apache Flink architecture overview
A Chinese deep dive into Flink's runtime topology, component responsibilities, and execution path.
- Reading
- 30 min
- Track
- Flink internals
- Source
- Chinese source notes
The source notes for this track are currently maintained in Chinese.
预计阅读时间: 30 分钟
前置阅读: 无
下一次阅读: doc-01(Job 提交调度) 或 doc-02(流数据链路)
1. Flink 是什么
Apache Flink 是流优先(streaming-first)的分布式计算引擎。核心定位:有状态的流处理,批处理是流处理的特殊形式(有界流)。
批处理 = 有界流(Bounded Stream)
流处理 = 无界流(Unbounded Stream)
2. 核心架构 — 进程拓扑

核心组件职责
| 组件 | 进程 | 类 | 核心职责 |
|---|
| Dispatcher | JM | Dispatcher.java | REST API + JobGraph 提交 → 启动 JobMaster |
| JobMaster | JM | JobMaster.java | JobGraph→ExecutionGraph + Scheduler + CheckpointCoordinator + Failover |
| ResourceManager | JM | ResourceManager.java | 管理 TM Slot, 向 YARN/K8s 申请 Container/Pod |
| Task | TM | Task.java | 单个 Task 的生命周期管理 |
| StreamTask | TM | StreamTask.java | processInput() 循环—Mailbox 模型 |
3. 完整调用链: 从 JAR 到运行
flink run job.jar
→ CliFrontend.run()
→ StreamGraphGenerator.generate() // DataStream API → StreamGraph
→ StreamingJobGraphGenerator.createJobGraph() // StreamGraph → JobGraph
→ Dispatcher.submitJob(JobGraph)
→ persistAndRunJob() // 持久化到 HA + 启动 JobMaster
→ JobMaster.onStart()
→ ExecutionGraphBuilder.buildGraph() // JobGraph → ExecutionGraph
→ DefaultScheduler.startSchedulingInternal()
→ SchedulingStrategy.startScheduling()
→ allocateSlotsAndDeploy(vertices)
→ RM.declareRequiredResources()
→ TM.startTask()
→ Task.doRun()
→ StreamTask.invoke()
→ processInput() 循环
4. 核心设计决策速查
| 决策 | 选择 | 核心原因 |
|---|
| 计算模型 | 流优先 | 批处理 = 有界流的特例 |
| 一致性 | Chandy-Lamport + Barrier | 分布式快照经典算法 |
| 反压 | Credit-based | 接收端驱动,无阻塞 |
| 状态 | StateBackend 抽象 | RocksDB (TB级) / Heap (100MB内) |
| SQL | Apache Calcite | 业界标准优化框架 |
| 内存 | Managed Memory + Network Memory | Off-heap 管理,避免 GC |
5. 代码库地图
Runtime 核心
| 组件 | 路径 | 入口类 |
|---|
| Dispatcher | flink-runtime/.../dispatcher/ | Dispatcher.java::submitJob() |
| JobMaster | flink-runtime/.../jobmaster/ | JobMaster.java::onStart() |
| ExecutionGraph | flink-runtime/.../executiongraph/ | DefaultExecutionGraph.java::transitionState() |
| Scheduler | flink-runtime/.../scheduler/ | DefaultScheduler.java::allocateSlotsAndDeploy() |
| ResourceManager | flink-runtime/.../resourcemanager/ | ResourceManager.java::registerTaskManager() |
| Task | flink-runtime/.../taskmanager/ | Task.java::doRun() |
Streaming 核心
| 组件 | 路径 | 入口类 |
|---|
| StreamTask | flink-streaming-java/.../tasks/ | StreamTask.java::processInput() |
| Operator Chain | flink-optimizer/.../ | StreamingJobGraphGenerator.java::isChainable() |
| WindowOperator | flink-streaming-java/.../windowing/ | WindowOperator.java::processElement() |
| Source (FLIP-27) | flink-core/.../connector/ | Source.java::createReader() |
Checkpoint & State
| 组件 | 路径 | 入口类 |
|---|
| CheckpointCoordinator | flink-runtime/.../checkpoint/ | CheckpointCoordinator.java::triggerCheckpoint() |
| RocksDB Backend | flink-state-backends/flink-statebackend-rocksdb/ | RocksDBKeyedStateBackend.java |
| CheckpointBarrierHandler | flink-streaming-java/.../io/checkpointing/ | SingleCheckpointBarrierHandler.java::processBarrier() |
网络栈
| 组件 | 路径 | 入口类 |
|---|
| RecordWriter | flink-runtime/.../network/api/writer/ | RecordWriter.java::emit() |
| ResultPartition | flink-runtime/.../network/partition/ | ResultPartition.java |
| Credit Flow Control | flink-runtime/.../network/netty/ | CreditBasedPartitionRequestClientHandler.java |
下一步