Flink 源码与运行原理

流式计算 / 源码阅读 / LESSON 00

Apache Flink 架构概览(深度篇)

先建立 JobManager、TaskManager、调度器、状态与网络栈之间的整体运行时地图。

阅读时间
30 分钟
学习路径
Flink 源码与运行原理
内容来源
Flink 深度笔记

预计阅读时间: 30 分钟 前置阅读: 无 下一次阅读: doc-01(Job 提交调度) 或 doc-02(流数据链路)


1. Flink 是什么

Apache Flink 是流优先(streaming-first)的分布式计算引擎。核心定位:有状态的流处理,批处理是流处理的特殊形式(有界流)。

批处理 = 有界流(Bounded Stream)
流处理 = 无界流(Unbounded Stream)

2. 核心架构 — 进程拓扑

Apache Flink 架构概览(深度篇) 图 01

核心组件职责

组件进程核心职责
DispatcherJMDispatcher.javaREST API + JobGraph 提交 → 启动 JobMaster
JobMasterJMJobMaster.javaJobGraph→ExecutionGraph + Scheduler + CheckpointCoordinator + Failover
ResourceManagerJMResourceManager.java管理 TM Slot, 向 YARN/K8s 申请 Container/Pod
TaskTMTask.java单个 Task 的生命周期管理
StreamTaskTMStreamTask.javaprocessInput() 循环—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内)
SQLApache Calcite业界标准优化框架
内存Managed Memory + Network MemoryOff-heap 管理,避免 GC

5. 代码库地图

Runtime 核心

组件路径入口类
Dispatcherflink-runtime/.../dispatcher/Dispatcher.java::submitJob()
JobMasterflink-runtime/.../jobmaster/JobMaster.java::onStart()
ExecutionGraphflink-runtime/.../executiongraph/DefaultExecutionGraph.java::transitionState()
Schedulerflink-runtime/.../scheduler/DefaultScheduler.java::allocateSlotsAndDeploy()
ResourceManagerflink-runtime/.../resourcemanager/ResourceManager.java::registerTaskManager()
Taskflink-runtime/.../taskmanager/Task.java::doRun()

Streaming 核心

组件路径入口类
StreamTaskflink-streaming-java/.../tasks/StreamTask.java::processInput()
Operator Chainflink-optimizer/.../StreamingJobGraphGenerator.java::isChainable()
WindowOperatorflink-streaming-java/.../windowing/WindowOperator.java::processElement()
Source (FLIP-27)flink-core/.../connector/Source.java::createReader()

Checkpoint & State

组件路径入口类
CheckpointCoordinatorflink-runtime/.../checkpoint/CheckpointCoordinator.java::triggerCheckpoint()
RocksDB Backendflink-state-backends/flink-statebackend-rocksdb/RocksDBKeyedStateBackend.java
CheckpointBarrierHandlerflink-streaming-java/.../io/checkpointing/SingleCheckpointBarrierHandler.java::processBarrier()

网络栈

组件路径入口类
RecordWriterflink-runtime/.../network/api/writer/RecordWriter.java::emit()
ResultPartitionflink-runtime/.../network/partition/ResultPartition.java
Credit Flow Controlflink-runtime/.../network/netty/CreditBasedPartitionRequestClientHandler.java

下一步