一文吃透 Flink 作业的“逻辑-物理”全链路
·
一、三层视角:从应用到作业到任务
Flink Application(你的程序入口)
你的 main() 构建管道(DataStream / Table / SQL),提交后形成逻辑图(Logical Graph):有向图,节点是 Operator(算子),边是数据流/数据集。
- Transformation:API 概念(如
map/keyBy/window),多由具体算子实现。 - Record:在流/集中的最小元素;函数与算子以 Record 为 I/O。
提交后:JobManager 接管
JobManager 是集群的“总指挥”,包含三组件:
- ResourceManager(RM):调配 Task Slot(资源调度单位)。
- Dispatcher:REST 提交入口 & Web UI;每个作业拉起一个 JobMaster。
- JobMaster(JM):单作业的“项目经理”,负责任务编排、故障恢复、Checkpoint 协调等。
下钻到执行:物理图(Physical Graph)
物理图是逻辑图的“可执行翻译版”:
- Task(任务):物理图的节点,是运行时最小工作单元。
- Sub-Task(子任务):同一算子/算子链的并行实例,各自处理数据的分区(Partition)。
- Operator Chain(算子链):两个或更多相邻算子无重分区直接相连,链内不经序列化与网络栈,大幅降低开销。
**TaskManager(TM)**是 Worker 进程,承载多个 Slot;Task 被调度到 TM 执行,TM 之间负责数据交换。
二、执行模式与集群形态:两两选对
Execution Mode(运行模式)
- STREAMING:默认流式,支持事件时间、水位线、状态/Checkpoint。
- BATCH:把有界输入当作“有界流”跑,故障恢复走“全量回放”,不启用 Checkpoint。
Cluster:Session vs Application
- Flink Session Cluster:长跑集群,可收多作业,生命周期不绑定作业;资源共享、启动快,但隔离性一般。
- Flink Application Cluster:一应用一集群,
main()在集群端运行;RM/Dispatcher 专属该应用,隔离强,云原生部署友好。
经验选型:
交互式/短作业多 → Session;
生产长驻/强隔离/易治理 → Application。
三、状态、后端与结果:长寿命与可演进
Managed State(托管状态)
已向框架注册的应用状态,Flink 负责**持久化、伸缩(rescaling)**等;与 Checkpoint/Savepoint 协作保障一致性与演进。
State Backend(状态后端)
决定状态存储位置与快照实现:
- JVM Heap(HashMap)→ 小状态、超低延迟;
- 嵌入式 RocksDB → 大状态、稳定、支持增量快照。
最小配置示例(Java):
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(10_000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// 建议:生产用持久化对象存储
env.getCheckpointConfig().setCheckpointStorage("s3://flink/ckpt");
// RocksDB 大状态
// env.setStateBackend(new EmbeddedRocksDBStateBackend());
JobResultStore
将已全局终止(完成/取消/失败)的作业结果持久化到文件系统,便于在 HA 场景下判定是否需要恢复,结果可“比作业更长寿”。
四、UID 与 UID hash:演进“锚点”的关键
- UID:算子的稳定唯一标识(可手动指定,也可由作业结构推导)。
- UID hash(Operator ID / Vertex ID):运行时标识,来源于 UID,出现在日志/REST/指标中,Savepoint 就靠它定位算子。
上线强烈建议:为有状态算子显式设置 UID,避免代码重构导致 ID 漂移,无法从保存点恢复。
DataStream<Event> ds = ...
ds.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(10)))
.reduce(new MyReduce())
.uid("user_10m_window_reduce") // 显式 UID
.name("10m-window-reduce");
排错提示:
- Savepoint 恢复报 “can’t map state to operator” → 检查 UID 是否变动。
- 多环境构建差异 → 固化
uid(),避免因优化/重排引发 hash 改变。
五、从逻辑到物理:链路是怎样落地的?

Source → Map → Filter可链化,落到 同一个 Task(同线程);keyBy/partition发生 repartition,生成新的下游 Task;- 下游 Task 的 Sub-Task 数量 = 并行度,各自消费不同 Partition。
六、Table Program:关系式管道与 DataStream 的桥
Table Program 是 Table API/SQL 声明的流水线。优势:语义清晰、优化器加持、上线快;可以与 DataStream 互转,混合编程“声明式为主、关键点下沉”。
示例(SQL 窗口 TVF):
CREATE TABLE events (
user_id STRING,
value BIGINT,
ts TIMESTAMP_LTZ(3),
WATERMARK FOR ts AS ts - INTERVAL '2' MINUTE
) WITH (...);
INSERT INTO sink
SELECT user_id, WINDOW_START, WINDOW_END, COUNT(*) AS cnt
FROM TABLE(
TUMBLE(TABLE events, DESCRIPTOR(ts), INTERVAL '10' MINUTES)
)
GROUP BY user_id, WINDOW_START, WINDOW_END;
七、部署与运行:Session/Application 两套食谱
Session Cluster(多作业共享)
- 先起集群,后多次
flink run; - 优点:冷启动快;缺点:一个 TM 故障可能牵连多个作业。
Application Cluster(强隔离推荐)
- 打包 JAR 即部署,入口在集群端调用
main()产出 JobGraph; - RM/Dispatcher 专属该应用,治理/升级边界清晰。
K8s Application 模式(示意):
flink run-application \
-t kubernetes-application \
-Dkubernetes.cluster-id=payment-app \
-Dtaskmanager.numberOfTaskSlots=2 \
local:///opt/flink/usrlib/job.jar
八、工程落地清单
- UID 固化:所有有状态算子显式
uid()。 - 状态后端:小状态 Heap;大状态 RocksDB + 增量快照;配置 Checkpoint 存储到可靠介质。
- Operator Chain:确保可链段链化(除非为隔离/排错手动断链)。
- 并行与分区:并行度 = 瓶颈算子并行;热点 Key 需加盐/预聚合。
- 执行模式:流式默认 STREAMING;有界数据可用 BATCH。
- 集群形态:交互短作业选 Session;长驻生产选 Application。
- HA 与恢复:开启 HA,配置 JobResultStore;定期 Savepoint 做版本演进锚点。
- 观测与告警:Backpressure、Checkpoint Duration、Busy Time、State Size、GC、Watermark Lag。
九、常见问题速解
- 保存点恢复失败:UID 改动或算子拓扑重排 → 回退代码或使用
uid()锚定。 - 吞吐低/延迟高:算子未链、slot 过度隔离、下游 Sink 阻塞 → 检查链化、Slot Sharing 与背压。
- TM 内存溢出:RocksDB 内外内存未限、slot 太多 → 收紧 block cache / write buffer,合理规划 TM 内存与 slots。
- 会话集群“牵一发而动全身”:关键作业迁往 Application Cluster,提升故障隔离性。
更多推荐


所有评论(0)