一、三层视角:从应用到作业到任务

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

八、工程落地清单

  1. UID 固化:所有有状态算子显式 uid()
  2. 状态后端:小状态 Heap;大状态 RocksDB + 增量快照;配置 Checkpoint 存储到可靠介质。
  3. Operator Chain:确保可链段链化(除非为隔离/排错手动断链)。
  4. 并行与分区:并行度 = 瓶颈算子并行;热点 Key 需加盐/预聚合。
  5. 执行模式:流式默认 STREAMING;有界数据可用 BATCH。
  6. 集群形态:交互短作业选 Session;长驻生产选 Application。
  7. HA 与恢复:开启 HA,配置 JobResultStore;定期 Savepoint 做版本演进锚点。
  8. 观测与告警: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,提升故障隔离性。
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐