精选·Flink

怎么理解 Flink 的流批一体

91学AI·2026/7/13·7 阅读

考察点

考察你对 Flink 设计哲学的理解深度。只会说「Flink 既能跑流又能跑批」是不够的,面试官想听的是为什么流批一体在 Flink 里是成立的、底层怎么实现的。追问会往有界流和无界流的差异、Table API 的统一、以及流批一体在数仓里的价值走。

参考答案

核心观点:批是流的特例

Flink 的世界观和 Spark 正好反过来。Spark 是批引擎模拟流(micro-batch),Flink 认为一切数据都是流:实时数据是无界流(unbounded stream),历史数据是有界流(bounded stream),批处理就是处理有界流这个特例。既然批是流的子集,那一个纯流引擎自然可以统一两者,而不是两套引擎拼凑。

这个观点不是口号,它决定了架构:Flink 的 DataStream API、Checkpoint 机制、状态管理全都围绕流设计,批处理复用同一套运行时(streaming runtime),只是数据源有界、作业会自然结束。

底层怎么统一

有界流和无界流在运行时上的差异主要体现在三点:

  • 调度方式。无界流所有算子同时部署、一直跑;有界流可以分阶段调度(像传统批处理的 stage),前序算子跑完落盘再跑下游,可以做更激进的优化。
  • 容错。无界流靠 Checkpoint 恢复;有界流中间结果可以持久化到磁盘(blocking shuffle),失败后只需重跑失败的阶段,不用从头回放。
  • 结束语义。有界流处理到输入末尾自然结束,Watermark 可以到正无穷,窗口全部触发。

Flink 1.12 之后 DataStream API 层面就支持了批执行模式(StreamExecutionEnvironmentRuntimeExecutionMode.BATCH),一套代码两种跑法;到 1.15 前后流批一体的 Shuffle、调度都成熟了。

Table API / SQL 层的统一

对业务开发来说,流批一体更多发生在 Table API 和 SQL 层:同一段 SQL,表来自 Kafka 就是流查询,来自 Hive/Iceberg 就是批查询,语义完全一致(窗口、Join、聚合的语义对齐)。这也是 Flink SQL 在实时数仓里能替代一部分离线 Hive 任务的原因:口径统一,离线 T+1 的结果和实时结果能对得上,Lambda 架构可以向 Kappa 演进。

工程价值

统一带来的是:一套代码维护、一套语义口径、一套运维体系。但也别吹过头——流批一体解决的是引擎和 API 统一,存储层(消息队列 vs 湖仓表)还是得选型,真正的「一份数据两处用」要靠 Paimon、Iceberg 这类流批一体的表存储配合,这通常是个好的延展点。

可能的追问

  • Spark Structured Streaming 也是流批统一吗?——它是批内核模拟流,微批模型有固定延迟下限,事件时间和状态的表达力弱于 Flink 的纯流模型;「流批一体」的成色不一样。
  • 批模式为什么性能可能更好?——有界输入让优化器能拿到完整统计信息做 cost-based 优化,可以选 sort-merge join、自适应并行度等策略,这是无界流做不到的。
  • 你们生产上怎么用流批一体的?——常见答案是:实时链路用流模式跑 Kafka,同一套 SQL 在湖表上跑批做数据回刷和口径校验。

评论 (0)

暂无评论,快来抢沙发吧!

91学AI

© 2026 91学AI · 按岗位学 AI 与大数据. All rights reserved.