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. 核心架构 — 进程拓扑

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

下一步