0%

Flink 流批一体与状态管理

Flink 流批一体与状态管理

一句话定位:Flink 的杀手锏是”一套引擎跑流和批”(流批一体),以及可扩缩、可容错的状态管理(State + Checkpoint);二者共同支撑 Exactly-Once。面试重点在 State Backend 选型、Barrier 对齐、以及 Exactly-Once 是怎么串起来的。

1. 流批一体(Unified Engine)

  • Flink 认为批是流的特例(有界流):同一套 DataStream API / Table-SQL API,通过 ExecutionMode.BATCH 即可跑有界数据,无需维护两套代码。
  • 统一 Runtime:调度、容错、shuffle、状态后端对批流一致;降低了”流一套、批一套”的维护成本,也保证了语义一致。
  • 面试话术:把”流批一体”讲成”统一 API + 统一 Runtime + 有界即无界特例”,比单纯说”支持批处理”更有深度。

2. 状态管理:State Backend 选型

Backend 存储位置 优势 适用
HashMapStateBackend JVM 堆 读写最快、无需序列化 小状态、低延迟
EmbeddedRocksDBStateBackend 本地磁盘 + 堆外 支持超大状态、增量 Checkpoint 大状态、长窗口、高吞吐
  • 选型看状态规模延迟预算:状态放得下堆就用 HashMap(快),放不下或要做增量 checkpoint 就用 RocksDB(慢但可扩展)。
  • 状态类型:Keyed State(按 key 分区,随 key 路由)/ Operator State(算子级,如 source offset)。

3. Checkpoint 的 Barrier 对齐机制

  • Barrier 由 Source 按固定间隔注入数据流,随数据向下游流动,标记”属于第 n 个 checkpoint 的数据边界”。
  • 对齐(Alignment):算子收到某个输入通道的 barrier-n 后,该通道后续数据被缓存阻塞,直到所有输入通道都收到 barrier-n,才做状态快照(snapshot)并向下游广播 barrier-n。
  • 对齐的意义:保证 Exactly-Once——barrier 之后的数据绝不会混入本次 checkpoint 的状态。
  • 非对齐 Checkpoint(Unaligned):把正在传输中的缓冲也纳入快照,避免反压下 barrier 被”堵”在通道里导致 checkpoint 超时,降低尾延迟,但快照更大。

4. Exactly-Once 如何达成

端到端 Exactly-Once = 可重放 Source + Checkpoint 状态持久化 + 事务/幂等 Sink

  1. Source 记录位点(Kafka offset / 文件偏移),失败可从 checkpoint 重放;
  2. 基于 Chandy-Lamport 分布式快照算法,barrier 对齐保证各算子状态对应同一逻辑时间点;
  3. Sink 用两阶段提交(2PC,如 Flink-Kafka 事务)幂等写入,使”状态提交”与”外部输出”原子化。

参考:Apache Flink 官方文档 State Backends / Checkpointing;论文《Lightweight Asynchronous Snapshots for Distributed Dataflows》(Chandy-Lamport)