Flink(一):核心概念、运行时架构与部署模式
Flink(一):核心概念、运行时架构与部署模式
导语:Flink 面试的地基是「流的世界观 + 运行时三件套 + 三层图」。本篇先讲清有界流/无界流与流批一体的本质,再拆解 Client / JobManager / TaskManager / Slot 的协作关系与「一 Job 一 JobMaster」的资源模型,然后梳理 StreamGraph → JobGraph → ExecutionGraph 的演进,最后落到四种部署模式与 Flink 2.0 的架构变化,共 14 题。
一、核心概念与定位
1. 为什么说 Flink 是「真正的流处理」?与 Spark Streaming、Storm 有什么区别?
答: 关键区别在于计算模型,而不是 API 风格。
| 维度 | Flink | Spark Streaming | Storm |
|---|---|---|---|
| 计算模型 | 原生流(逐条处理,事件驱动) | 微批(把流切成小批次 RDD) | 原生流(逐条) |
| 延迟 | 毫秒级 | 秒级(批间隔决定) | 毫秒级,最低 |
| 吞吐 | 高 | 高 | 相对低 |
| Exactly-Once | 原生支持(Checkpoint + 两阶段提交) | 支持(需幂等/事务 Sink) | At-Least-Once(Trident 才能 EO) |
| 时间语义 | Event Time + Watermark 完整支持 | 有限支持 | 弱(主要 Processing Time) |
| 状态管理 | 一等的 Keyed State / Operator State + StateBackend | 靠 RDD 血统重算或外部存储 | 需自己接外部存储 |
| 背压 | Credit-based 流控,天然背压 | 微批天然背压 | 反压机制弱,易 OOM |
| 流批一体 | 同一套 API + 同一套运行时 | 流批 API 分离(Structured Streaming 稍好) | 无 |
一句话总结:
- Storm:延迟最低,但状态、Exactly-Once、Event Time 都弱,现在基本被替代;
- Spark Streaming:微批模型导致延迟下限受批间隔限制,且「流」是批的模拟;
- Flink:把「有状态计算」和「事件时间」做成第一公民,并用 Checkpoint 统一解决容错与 Exactly-Once,是当前实时计算的事实标准。
面试常追问:「微批 vs 原生流」的真实差异在哪? 答:微批的延迟下限 = 批间隔,且窗口/定时语义要靠批边界近似;原生流可以做到逐事件触发,并且Checkpoint 是异步的、与数据处理解耦,不需要停流等批结束。
2. 有界流和无界流是什么?流批一体怎么理解?
答:
- 无界流(Unbounded):有开始、没有结束,数据持续产生(如埋点日志、订单变更),必须持续处理,这是流处理的主战场;
- 有界流(Bounded):有开始、有结束,即批数据(如一天的离线文件)。
Flink 的核心设计哲学是「批是流的一个特例」——批就是「有界流」,因此可以用同一套算子、同一套运行时处理两者,只是执行策略不同:
| 执行模式 | 触发方式 | 关键差异 |
|---|---|---|
| STREAMING | 常驻运行 | 逐条处理;Checkpoint 保证 EO;Watermark 驱动窗口 |
| BATCH(原 Dataset/DataSet 语义) | 有界输入,处理完即退出 | 支持阻塞式 Shuffle(Blocking Shuffle,可落盘+可排序)、按 key 排序、更高效的批算子(如 sort-merge join)、不需要 Checkpoint 做容错(失败重算即可) |
// 同一段 DataStream 逻辑,执行模式由 env 决定
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setRuntimeMode(RuntimeExecutionMode.BATCH); // 或 STREAMING / AUTOMATIC版本演进(必知,容易体现深度):
- Flink 1.12 起 DataStream API 统一批流(FLIP-134),
DataSet API开始被弃用; - Flink 1.14/1.15 完善 BATCH 模式下的 Shuffle、排序与流水线阻塞;
- Flink 2.0 已正式移除
DataSet API,批处理统一走 DataStream/Table API 的 BATCH 模式; - 表层的统一是 Table API & SQL:同一段 SQL 在
STREAMING与BATCH模式下都能执行,只是语义不同(流模式是持续查询,批模式是一次性查询)。
3. Flink 的四大核心支柱(能力全景)是什么?
答: 可以概括为四件事,面试时先抛出这四个关键词,再展开细节,会显得结构清晰:
| 支柱 | 解决的问题 | 面试关键词 |
|---|---|---|
| Event Time + Watermark | 数据乱序到达,凭什么断定「窗口可以计算了」 | Watermark、allowedLateness、侧输出 |
| State(状态) | 跨事件、跨批次的上下文(累加值、去重集合) | Keyed/Operator State、StateBackend、State TTL |
| Checkpoint / Savepoint | 故障后状态不丢、Exactly-Once | Barrier 对齐、异步快照、两阶段提交 |
| 反压与内存管理 | 上下游速度不匹配、大状态不 OOM | Credit-based 流控、Managed Memory、Buffer Pool |
这四点是分工的:时间语义解决「何时算」,状态解决「算什么」,Checkpoint 解决「算错了怎么办」,反压解决「算不过来怎么办」。
二、运行时架构
4. 描述一下 Flink 的运行时架构(四个角色分别干什么)?
答:
| 角色 | 职责 | 关键点 |
|---|---|---|
| Client | 构建 StreamGraph、转换为 JobGraph、提交后在 Application 模式下即为 main() 进程 | 客户端不是运行时的一部分,提交后可退出(Session/Per-Job) |
| JobManager | 集群的「主」,又拆成 Dispatcher + ResourceManager + JobMaster 三个组件 | 高可用(HA)需要 ZK/K8s 选主 + 元数据持久化 |
| TaskManager | 集群的「从」,真正跑 Task 的 JVM 进程,一个 TM = 若干 Task Slot | 内存是预分配池化的(Managed Memory / Network Buffer) |
| Slot | TaskManager 上资源子集的抽象(不是线程,也不是固定核数),是调度的最小单位 | 默认一个 slot ≈ 一个 CPU 核(taskmanager.numberOfTaskSlots),实际按内存/核数配置 |
三个 Master 组件的分工(高频追问):
- Dispatcher:接收 Client 提交的作业、提供 Web UI 入口、为每个作业启动一个 JobMaster(Application 模式还会顺便启动
main()所在集群); - ResourceManager:只管 TaskManager 的 Slot(不是 YARN/K8s 的 RM),负责 Slot 的申请、分配与回收,并向外部资源管理器(YARN / K8s / Standalone)申请/释放容器;
- JobMaster:一个作业一个,负责把
JobGraph转成ExecutionGraph、申请 Slot、调度 Task、协调 Checkpoint、处理作业状态变更。
记忆口诀:Dispatcher 管「作业入口」,ResourceManager 管「Slot 资源」,JobMaster 管「这一个作业的执行」。一个 JobManager 进程里可以有多个 JobMaster(Session 模式),这就是「一个集群跑多个作业」的实现方式。
5. Slot、并行度、Slot Sharing 三者是什么关系?
答: 先分清三个概念:
| 概念 | 含义 | 谁来配 |
|---|---|---|
| 并行度(Parallelism) | 一个算子被拆成多少个子任务并行执行 | env.setParallelism() / 算子级 setParallelism / SQL table.exec / 提交参数 |
| Slot | TaskManager 提供的资源子集,用来放置 Task | taskmanager.numberOfTaskSlots |
| Slot Sharing Group | 允许共享同一 Slot 的算子集合的定义(默认所有算子属于 default 组) | 算子级 slotSharingGroup("g") |
Slot Sharing(槽位共享)的作用:默认情况下,同一作业的所有算子共享同一个 Slot(属于默认共享组)。这样做的收益是:
- 一个 Slot 就能装下整条链,不用为每个算子单独申请资源,作业所需 Slot 数 = 最大并行度(而不是算子数 × 并行度);
- 降低资源碎片:如果一个作业只在部分算子上有热点,共享后能把 CPU 利用起来(总核数不用按算子逐个配满);
- 减少网络传输:共享 Slot 的算子之间是同线程内存传递(见第 6 题)。
关键约束:
- 同一 Slot 内的并发子任务数由最"重"的算子决定——所以「Slot 数够了但 Task 抢不到资源」通常是因为并行度配置不一致;
- 一个 Slot 内的 Task 数量理论上=「Slot Sharing Group 中并行度最大的算子值」;
- 并行度 > Slot 总数时 Task 会排队等待,日志表现为
insufficient number of slots; slotSharingGroup隔离:把重算子和轻算子分开能避免互相影响,但跨共享组的传输必须走网络。
高频易错点:Slot 默认不等于 CPU 核的硬隔离,它是资源池里的逻辑份额;真正的资源隔离要靠配置
taskmanager.numberOfTaskSlots = CPU 核数,让每个 Slot 大致对应一个核,否则同一 TM 内多个 Slot 会争抢 CPU。
6. 算子链(Operator Chain)是什么?什么条件下算子能链在一起?
答: 算子链把多个算子合并成一个 OperatorChain,在同一个线程里依次执行,从而把数据传输从「网络 + 序列化」变成「对象引用直传」,大幅降低开销。
能链在一起的条件(两个必须同时满足):
| 条件 | 说明 |
|---|---|
| 上下游并行度相同 | 并行度不同则子任务无法一一对应,必须走网络 shuffle |
| 同一 Slot Sharing Group | 不在同一共享组则不可能同 Slot,也就不能同线程 |
不能链 / 需要打断的情况:keyBy、rebalance、broadcast、rescale 等需要重分区的算子;显式调用 startNewChain()、disableChaining()、slotSharingGroup("x") 也会打断链。
优缺点与调优建议:
- 优点:减少线程切换与序列化开销、降低延迟、提升吞吐;
- 缺点:链太长会放大反压传导,且一个子任务出问题会拖累整条链;链内算子无法单独设并行度;
- 调优:
- 想「让某些算子单独出来」→
disableChaining()/startNewChain(); - 想「同一算子的不同子任务不共享 Slot 的本地状态」→ 用
slotSharingGroup分开; - Web UI 上会显示
OperatorChain与「chain 中的算子列表」,是排查性能问题的第一手信息。
- 想「让某些算子单独出来」→
顺带一个常见面试问法:「Flink 为什么比 Spark 快?」 除了原生流,算子链 + 内存级 Shuffle + 序列化位置优化(把对象序列化下推到需要网络传输的那一步) 都是关键原因。
7. StreamGraph、JobGraph、ExecutionGraph 分别是什么?
答: Flink 从代码到执行要经过四层图,每一层解决不同问题:
| 层 | 谁构建 | 作用 | 关键概念 |
|---|---|---|---|
| StreamGraph | Client | 用户代码的逻辑拓扑 | StreamNode / StreamEdge;边上有 Partitioner(Forward/Hash/Rebalance…) |
| JobGraph | Client | 提交与优化的结果 | 算子链合并、JobVertex / IntermediateDataSet、执行配置 |
| ExecutionGraph | JobMaster | 并行化的执行计划 | ExecutionJobVertex → ExecutionVertex(每个并行子任务);Execution 可重试 |
| 物理执行图 | TaskManager | 真正跑起来的 Task | 每个 Task = 一个 OperatorChain 的 subtask |
为什么面试要问这个:它把「客户端优化(链化)→ 调度并行化(展开)→ 运行时容错(Execution 重试)」三段职责讲清楚了。比如:
- 算子链在 JobGraph 阶段就已确定,所以运行期无法改变链结构(改链只能改代码重新提交);
- 故障恢复的粒度是
Execution(一个并行子任务),而不是整个算子——ExecutionGraph 里的ExecutionVertex才是重试单位; - 动态扩缩容(Adaptive Scheduler / Reactive Mode) 改的就是
ExecutionGraph的并行度,需要先做 Savepoint 或依赖AdaptiveScheduler的空闲 Slot 重调度。
三、部署与版本演进
8. 四种部署模式(Session / Per-Job / Application)有什么区别?
答: 区别在「集群的创建时机、作业隔离度、谁来跑 main()」:
| 模式 | 集群启动时机 | main() 在哪执行 | 资源隔离 | 适用场景 |
|---|---|---|---|---|
| Session | 先起集群,后提交作业 | Client(本地) | 无隔离,所有作业共享 JobManager 与 Slot | 小作业多、快速提交、开发调试 |
| Per-Job | 一个作业起一个集群,作业结束集群销毁 | Client | 完全隔离 | 生产长作业(1.15 起已废弃) |
| Application | 一个应用(可含多作业)起一个集群 | 集群内(JobManager 侧) | 隔离到「应用」粒度 | 生产推荐,尤其 K8s |
| Session on YARN / K8s | 先起 Session 集群 | Client | 无(共享) | 交互式、低成本试跑 |
Session 模式的两个典型坑:
- Client 承担
main()与依赖打包:main()在客户端生成 JobGraph 后再提交,依赖 jar 冲突/连接被断都会导致提交失败;如果 Client 进程退出,Session 模式的作业不受影响(已经提交完),但 Application 模式语义不同; - JobManager 成为单点瓶颈:所有作业共用一个 JM,一个作业把 JM 的 Slot 申请/Checkpoint 协调压满,其他作业一起抖动;且本地依赖冲突(classloader 隔离弱)容易互相影响。
为什么生产推荐 Application 模式:
main()在集群内执行,客户端只负责提交,不参与依赖下载与代码生成,天然适合容器化 / 多租户;- 集群的生命周期 = 应用生命周期,资源隔离与回收干净;
- K8s 下配合 Native Kubernetes 部署或 Flink Kubernetes Operator,能按作业粒度做弹性与升级。
补充:Per-Job 模式在 1.15 被标记废弃,官方建议用 Application 模式替代,因为它能做到同样的隔离度且不需要独立的客户端资源。
9. 常见的资源管理与部署方式有哪些?
答:
| 部署方式 | 说明 | 要点 |
|---|---|---|
| Standalone | 手工起 JobManager / TaskManager,Flink 自己管资源 | 简单;不支持弹性扩缩容,生产少用 |
| YARN | Flink 把 YARN 当外部 RM 申请容器 | 传统大数据栈常见;支持 Session/Per-Job/Application |
| Kubernetes | Native K8s(JobManager 直接调 K8s API)或 Standalone K8s 或 Flink Operator | 主流;可 reactive 模式自动扩缩容 |
| 云托管 | 阿里云实时计算 Flink 版、AWS 等 | 免运维,支持 Serverless |
HA(高可用)要点:
- JobManager HA:需要选主组件(ZK / K8s Lease)+ 元数据持久化(
high-availability.storageDir存JobGraph、Checkpoint 指针等);备 JM 接管后从最近一次 Checkpoint 恢复,因此 Checkpoint 的可靠存储是前提; - TaskManager 无常:TM 挂了由 JobMaster 感知 → 触发作业重启 / Region 重调度;
- JobManager 是单点写:即使 HA,元数据写入仍串行,大规模作业要注意 JM 的 CPU/GC。
K8s / 云原生的关注点(加分项):
- StatefulSet vs Deployment:TaskManager 用 Deployment 即可(无状态),
reactive模式下由 Flink 自己控制副本数; - 本地状态 + PVC:RocksDB 本地目录建议挂 SSD/PVC,配合本地恢复(Local Recovery)加快恢复;
- 镜像是版本契约:Flink 版本与作业 jar 必须与镜像一致,否则反序列化失败;
- 弹性扩缩容:
reactive模式自动跟随 Slot 变化重调度,但有状态扩缩容需要 Savepoint 或自适应调度器支持,不是无成本。
10. Flink 2.0 有哪些值得关注的变化?
答: 1.x 到 2.0 是架构级的升级,面试提到能加分,但要注意只讲确定的事实:
1)状态「存算分离」(Disaggregated State Storage)
Flink 1.x 的状态是存算一体的:状态落在 TaskManager 本地(堆或 RocksDB 本地盘),Checkpoint 再周期性上传远端。
- 引入 ForSt(Flink Remote State) 这一套存算分离状态后端(含
ForStDB实现)与异步状态 I/O,目标是解决「百 TB 级大状态下的 Checkpoint 慢、恢复慢、扩缩容抖」; - 1.x 的
HashMapStateBackend/EmbeddedRocksDBStateBackend仍保留(兼容),但大状态方向推荐存算分离路径。
2)其他重要变更(知道即可)
- 移除
DataSet API:批统一走 DataStream 的BATCH模式与 Table API; - 要求 Java 17+:JDK 8/11 不再支持,序列化与内存管理也随之调整;
- 新 Source / Sink API(FLIP-27 / FLIP-143)成为唯一形态,旧的
SourceFunction/SinkFunction被移除; - 物化表(Materialized Table):用声明式方式统一「流」与「批」的产出表,落地流批一体的湖仓场景;
- Serverless / 云原生增强:更好支持弹性伸缩与存算分离部署。
记忆点:Flink 2.0 = 存算分离(ForSt)+ 移除 DataSet + JDK 17 + 新 Source/Sink + 物化表。
四、数据流抽象与开发模型
11. DataStream 编程的基本结构和常用算子有哪些?
答: 一个 Flink 作业的骨架永远是「环境 → Source → Transformation → Sink → execute」:
public class Demo {
public static void main(String[] args) throws Exception {
// 1. 环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(4);
env.enableCheckpointing(10_000); // 容错前提
// 2. Source(新 API:Source → SourceFunction 已废弃)
DataStream<String> lines = env.fromSource(
KafkaSource.<String>builder()
.setBootstrapServers("kafka:9092")
.setTopics("order")
.setGroupId("flink-order")
.setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))
.setValueOnlyDeserializer(new SimpleStringSchema())
.build(),
WatermarkStrategy.noWatermarks(),
"kafka-source");
// 3. Transformation
DataStream<Order> orders = lines
.map(JsonUtils::parse) // 1:1
.filter(o -> o.getAmount() != null) // 1:1
.keyBy(Order::getUserId) // 重分区(hash)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))// 窗口
.aggregate(new SumAmount()); // 增量聚合
// 4. Sink
orders.sinkTo(KafkaSink.<Order>builder()
.setBootstrapServers("kafka:9092")
.setRecordSerializer(...)
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.build());
// 5. 触发
env.execute("order-window-app");
}
}常用算子按「是否触发重分区」分类(这是理解成本的关键):
| 类别 | 算子 | 是否 shuffle | 说明 |
|---|---|---|---|
| 单流转换 | map / flatMap / filter | 否(Forward) | 保持并行度,可链化 |
| 聚合(有状态) | keyBy + sum/max/min/reduce/aggregate | keyBy 是 Hash 重分区 | 增量聚合,状态可控 |
| 富函数 | ProcessFunction / KeyedProcessFunction | 视上游 | 可访问状态、Timer、侧输出,最灵活 |
| 窗口 | window / windowAll / countWindow | 沿用上游分区 | 见《Flink(二)》 |
| 连接/合并 | union(多流合并,类型一致) | 否 | 只是逻辑合并 |
connect(类型可不同) | 否 | 常配 CoProcessFunction | |
join / intervalJoin / coGroup | 是(keyBy) | 双流关联,有状态 | |
| 分流/重分区 | sideOutput(侧输出) | 否 | 分流推荐方式(split 已废弃) |
keyBy / rebalance / rescale / broadcast / forward / shuffle / partitionCustom | 是 | 见第 12 题 | |
| Sink | sinkTo / addSink | 视实现 | 新 API 用 sinkTo |
两个容易踩的实现坑:
keyBy之后不能再windowAll:windowAll要求并行度为 1(非 keyed 全量窗口),只能用在keyBy之前或parallelism=1的场景;ProcessFunction里的状态必须通过RuntimeContext获取(getRuntimeContext().getState(...)),且必须在open()方法里初始化,否则每次processElement都会重复创建/或在序列化阶段报错。
12. Flink 的 Partitioner(分区策略)有哪些?
答: Partitioner 决定上游子任务的输出发到下游哪个子任务(有些还会触发网络传输):
| 策略 | 语义 | 是否网络传输 | 典型用途 |
|---|---|---|---|
| Forward | 上下游 1:1(同索引子任务) | 否 | 并行度相同时的直通,可链化 |
Hash(keyBy) | hash(key) % 下游并行度 | 是 | 保证同一 key 落到同一子任务(有状态聚合的前提) |
| Rebalance | 轮询分发 | 是 | 打散热点、均衡负载 |
| Rescale | 本地轮询:每个上游只发到「对应的一组」下游 | 是 | 比 rebalance 更省网络(上下游按比例映射,只走本节点内) |
| Broadcast | 复制到所有下游子任务 | 是 | 小表广播(如配置、维表全量) |
| Global | 全部发到下游第一个子任务 | 是 | 少用,易成瓶颈 |
| Shuffle | 随机 | 是 | 随机打散 |
| Custom | 自定义(实现 Partitioner) | 视实现 | 特殊路由需求 |
要点:
keyBy用的是 hash 分区,所以「同一 key 一定在同一并行子任务」是有状态计算的正确性基础(否则累加值会被拆到不同子任务);rescale与rebalance的区别是高频考点:rebalance是所有上游把数据轮询发给所有下游(跨节点多),rescale是上游只发给「自己所在的那一组下游」(减少跨节点传输,但可能不均衡);- 实际数据倾斜时优先考虑:
rebalance/rescale打散,或两阶段聚合(加盐 + 去盐),见《Flink(四)》; keyBy的 key 必须实现hashCode()与equals()一致,否则会因「两个逻辑相等的 key 分到不同子任务」而产生业务错误的状态(不易发现)。
13. Flink 的编程模型里「UDF 序列化」是怎么回事?为什么要求类型信息?
答: Flink 需要把对象在网络传输、状态存储、Checkpoint 之间搬运,因此必须知道「怎么把对象变成字节」。
三条路径(性能从高到低):
- POJO 类型(推荐):类满足「public 类 + public 无参构造 + 所有字段 public 或带 getter/setter」,Flink 用
PojoSerializer按字段序列化,字段级高效、Schema 演进友好(可加字段); - 基础类型 / 数组:如
String、int、Row、Tuple,有专用序列化器,最快; - Kryo 兜底(尽量避免):不满足 POJO 条件的普通 Java 类会退回 Kryo,灵活但慢(反射+全字段),且类结构变化后无法兼容旧状态/Checkpoint。
实践建议:
| 建议 | 原因 |
|---|---|
优先定义 POJO 或用 Tuple/Row | 序列化性能差一个数量级;Kryo 容易在升级时炸 |
用 TypeInformation 显式声明 | 避免泛型擦除导致 Flink 退化成 Kryo;Lambda 中带泛型尤其容易 |
| 避免在状态里存大对象 | 状态会被反复序列化/反序列化,大对象会拖垮 Checkpoint 与吞吐 |
不要依赖 Kryo 的默认注册 | 未注册的类会有日志告警,且不同作业间不兼容 |
面试常见问法:「为什么我的作业日志里一直提示
GenericTypeInfo/ 建议注册 Kryo?」 答:类型信息丢失(多为 Lambda/泛型/私有字段),Flink 只能用 Kryo 兜底——功能能跑,但性能和兼容性都会打折。
14. 面试总结:一条 Flink 作业从提交到运行的完整链路
答: 把本篇串成一条线,面试时能一口气讲清楚,基本就过了「架构」这一关:
关键记忆点:
- 客户端与运行时不强绑定:Application 模式下
main()在集群内; - Slot 是调度单位,
并行度 = 最大算子并行度,Slot 数 = 并行度(默认共享组); - 算子链在 JobGraph 阶段定死,运行期改不了;
- 重试粒度是 Execution(子任务),不是整个算子;
- 一切容错都建立在「状态 + Checkpoint」上,所以下一篇才是真正的主战场。
下一篇:《Flink(二)》深入时间语义与窗口——三种时间语义的取舍、Watermark 的生成与传播、乱序与迟到数据的三种补救(allowedLateness / 侧输出 / 重算),以及窗口类型、分配器、触发器与窗口生命周期。
