聊到批处理,很多人第一反应就是“Spark的天下”。过去的几年里,我个人在多个数据平台上同时跑过Spark和Flink的离线任务,前期用Spark做数仓加工和批量ETL,后面为了湖仓一体架构,硬是把一部分批任务迁到了Flink上。做完这一轮迁移,一个很真实的感触是:Flink在流处理上确实能打,但批处理要想全面替代Spark,还差着一大截。这篇文章不聊流处理,专门聚焦批处理场景下,Flink相对Spark到底有哪些实质性的不足,哪些坑是我真实踩过的,以及在实际选型时该怎么权衡。
这可能是想对比两个引擎的人最想看到的内容。如果你是准备入行大数据、正在做技术选型、或者已经在用Flink做批任务但总觉得“哪里不对劲”,这篇内容会给你一个相对全面的参考。
1. “流批一体”的旗号下,Flink骨子里仍然偏向流
1.1 设计哲学的差异:一个是天生的批引擎,一个是穿上西装跑批的流引擎
我说句实在话,Flink最引以为傲的标签就是“流批一体”,官方文档里也一直在强调“用一套引擎搞定流和批”。但实际用下来你会发现,所谓的流批一体更像是“以流为底座,把批当成流的特殊情况”——这和Spark“以批为底座,把流当成微批”的路径,本质上就是两条完全不同的路线。
Spark从第一天开始就是为批量计算设计的。RDD、DAG、Stage、Shuffle,这些概念全部围绕“一批数据一次性处理完”来构建。哪怕后来的Structured Streaming用微批模型做流处理,它的执行引擎底层依然是批处理那一套,只是在时间维度上把数据切成了一段一段的小批量。这种设计的好处是,批处理这个主业永远是最稳、最成熟的。
Flink走的是另一条路。它的核心抽象是连续流和状态,DataStream、Window、Checkpoint、Event Time……这些概念统统围绕流来设计。批处理在Flink里被定义为“有界流”,也就是说,把数据当作一个有限的流来处理。理想很丰满,但现实是执行引擎里很多组件、参数默认值、资源模型,都是先为流设计的。真到了跑批任务时,你经常会感觉到一种“穿着西装去工地搬砖”的别扭。
1.2 三个看得见摸得着的后果
这种设计哲学上的差异,会直接在运行时有三个实际表现:
- 调度最小单位不同。Spark的调度单位是Stage内部的Task,每个Stage可以独立调度,边算边跑;Flink的Job图在启动时就把整条链路的算子全部铺开,所有算子共同组成一个执行拓扑,即使批处理场景下很多算子之间有必要的数据交换,它也是“整条流水线一起动”的思路。
- 资源生命周期不同。Spark每个Stage执行完,相关Executor资源就可以被释放或复用给其他Stage;Flink一旦启动,整个拓扑申请的资源就会一直占用到Job结束,在批任务长尾阶段尤其浪费。
- 容错粒度不同。Spark靠血统(Lineage)重算失败的分区,灵活轻量;Flink主要靠Checkpoint与状态恢复,这套机制在流任务里天经地义,但在批任务里很多时候就是额外开销。
别小看这三点,它们组合在一起,就是后面几章里Flink批处理各种劣势的总根源。
2. 优化器和SQL执行能力的差距,是批场景最直接的痛
2.1 Catalyst/Tungsten/AQE的成熟度,Flink目前没能完全追平
批处理场景和流处理场景最大的不同在于,批任务的数据规模通常要大好几个量级,查询复杂度也高得多。这种场景下,SQL优化器的能力几乎决定了你的任务能不能跑得动、跑得快不快。
Spark的SQL优化体系经历了十来年的大规模磨炼:Catalyst优化器负责逻辑计划和物理计划的优化,Tungsten把代码生成、内存管理做到了极致,自适应查询执行(AQE)在Spark 3.0正式落地之后,更是把动态合并Shuffle分区、自动处理数据倾斜、动态切换Join策略这些能力都变成了默认选项。
Flink这边的优化器发展就慢了不少。虽然Flink也有自己的查询优化器,1.18版本之后也在大改基于成本的优化(CBO)和新的计划框架,但和Spark的成熟度相比还是差着代际:
- 自适应能力不如Spark完整。Flink直到近几个版本才在做动态并行度调整这类自适应能力,而且应用范围和稳定性远不如Spark AQE。真实场景中,Spark遇到数据倾斜可以自动加盐、自动优化Join策略,Flink往往需要你手工去排查倾斜Key,手工加盐、手工改并行度。
- 统计信息收集和估算能力弱。Spark CBO会基于表统计信息、列统计信息做代价估算,选择更优的执行计划;Flink的优化器对统计信息的依赖和利用还比较有限,很多复杂SQL只能靠经验手工调整。
- 谓词下推、分区裁剪等优化规则的覆盖面不如Spark全面。在复杂的多表关联场景下,Flink有时会生成一些明显不够聪明的执行计划,比如Filter不下推、子查询处理不好,需要DBA介入改写SQL。
2.2 一个非常典型的场景:大规模JOIN和倾斜处理
我用一个真实场景来说话。有一次数仓里两张表关联,一张是大事实表,几十亿条记录,一张是维表,几千万条记录,关联条件存在明显热点Key。在Spark上跑,AQE会自动将Shuffle分区数量从2000降到合适水平,自动识别倾斜分区并进行拆分,整个过程几乎不用人工干预,任务稳定在十几分钟跑完。
同样的数据、同样的SQL,我在Flink上跑了一遍。结果不仅没有自动倾斜优化,Hash Join的并行度也没选好,导致个别子任务卡出了长尾,最终整个任务跑了将近四十分钟。我后来不得不手动找出热点Key、手动加盐打散,才把时间压回到二十分钟内。
这种体验非常能说明问题。批处理场景里,数据倾斜几乎是必然存在的,Spark已经把处理倾斜变成了“平台自动完成的事”,而Flink在这个方向上的工具化程度远远不够。对做数据开发的同学来说,这就意味着你需要花更多精力在业务逻辑之外的任务调优上。
2.3 从TPC-DS基准测试看行业差距
TPC-DS是最主流的批查询性能基准,包含大量复杂分析查询:多表关联、多层子查询、窗口函数、CUBE/ROLLUP应有尽有。Spark社区这么多年一直在围绕TPC-DS做持续的性能优化和计划修正,业界已经有大量基于Spark跑TPC-DS的公开优化实践。
Flink在批处理侧的TPC-DS表现,从产业界的反馈和公开测试来看,整体上仍然不如Spark,尤其在查询计划很复杂、关联层次很深的任务上,需要投入的调优成本更高。Flink社区也承认批处理优化器在过去几年里相对薄弱,目前的版本已经在补课,但补课的效果还需要时间检验。
2.4 开发体验:Spark的API全家桶 vs Flink批处理的“偏科生”
批处理开发的日常不光是写SQL。很多时候你会用到DataFrame的算子、窗口函数、UDF、COALESCE调整文件数,还要处理各种格式转换。
Spark在这方面的开发体验几乎没有短板:DataFrame API功能齐全,PySpark让Python用户也能丝滑操作分布式数据,Spark SQL支持inferSchema自动推断JSON文件结构,读一个JSON文件可以直接spark.read.json()一把梭。
Flink的Table API虽然功能上也在逐步补齐,但很多能力默认围绕流式语义展开。批处理场景下,API的表现力和周边库的丰富程度明显不如Spark。比如处理嵌套JSON时,Flink SQL往往需要预先定义完整的ROW类型,像Spark那样通过采样自动推断Schema的能力就弱了不少。这种开发体验上的“不爽”,平时不觉得,真正写任务的时候会一直膈应你。
3. 资源调度和执行引擎的“体质差异”:跑批时Flink吃亏在哪
3.1 调度粒度:点菜式调度和包场式调度
Spark的调度方式,我更喜欢叫它“点菜式调度”。每个Job按Stage拆分,前面Stage算完把结果落盘,后面Stage再启动。也就是说,任务不是一次性把所有资源吃满,而是按阶段推进、按需申请。哪怕是一个跑几十个Stage的大任务,也可以前一个Stage完事、释放资源,后一个Stage再继续申请。
Flink的调度方式就有点“包场式调度”的味道了。Job提交后,会一次性为整张执行图的所有算子申请Slot资源,整个执行拓扑建成后开始处理数据。批任务里如果某个算子只需要很短的时间处理,它占用的Slot也必须在整个Job生命周期内一直保留。
这在批处理里造成的问题很现实:集群里多跑几个Flink批任务时,资源会被预先占满,而实际并发计算的任务可能远没有那么多,整体资源利用率和吞吐并不理想。我维护的一个离线集群曾经是Flink批为主,一到调度高峰期就会出现任务排队等待资源,但集群整体的CPU利用率并不高——资源被Job级地预留,而不是被Task级地复用。
3.2 Shuffle机制:Spark久经考验,Flink还在路上
Shuffle是批处理的灵魂。两个引擎都需要在执行过程中实现数据重新分区,但机制差异很大。
Spark的sort-based shuffle经过了十几年的迭代,支持Map端输出合并、外部排序、大规模溢写,对于几十TB级的Shuffle场景有非常成熟的保障机制。Spark还能通过调整spark.sql.shuffle.partitions、spark.default.parallelism等参数灵活控制Shuffle并发。
Flink批处理在Shuffle这个环节的成熟度,说实话,还有不小的提升空间。Flink推行的“高效Shuffle”方案近年来有一些改进,但在大规模批处理Shuffle时,磁盘IO、网络传输和内存开销的控制能力依然不如Spark稳定。我在生产环境里遇到过,同样规模的Shuffle,Spark任务稳定跑完,Flink可能需要调大并行度和网络缓冲,否则容易出现反压和OOM。
这里补充一点,Flink的设计目标决定了它会优先保障流场景下的低延迟和状态一致性,Shuffle机制也要兼顾流批两套执行模式,因此在纯批处理这个场景下,它很难像Spark那样把Shuffle打磨到极致。
3.3 容错机制:血统重算和Checkpoint的成本差异
批处理任务最怕的就是运行到一半失败。两个引擎的容错思路完全不同。
Spark用血统(Lineage)机制,每个RDD分区都知道自己的数据是从哪里计算来的,某个分区丢失后只需要重新计算那一个分区。这种“哪里坏了补哪里”的方式,失败恢复成本很低,恢复速度也快,尤其适合有向无环图的批处理流程。
Flink的容错核心是Checkpoint。批处理模式下,虽然不需要像流处理那样持续做Checkpoint,但一旦任务失败,它依然依赖Checkpoint或作业重启机制来恢复。批任务的计算链路通常很长,中间有不少算子会持有状态,Checkpoint本身就需要序列化和持久化,这会产生额外的开销。状态大的场景下,恢复过程可能要重新加载大量状态,比Spark的血统重算慢不少。
说白了,Spark的容错更符合批处理“错哪补哪”的逻辑,Flink的容错更适应流处理“持续备份”的逻辑。拿到批场景里比,自然落下风。
3.4 动态资源分配:Spark自带头顶光环
Spark 1.2就引入了动态资源分配,任务量小的阶段可以自动释放Executor,集群可以把这些资源拿给其他任务用。Flink这边,批任务的资源动态调整能力一直相对薄弱。尽管Flink 1.17之后在自适应调度上有一些进展,比如根据数据量和吞吐自动调整并行度,但整体上还没有达到Spark dynamic allocation那种成熟的资源弹性水平。
我个人的体感是:在一个同时跑Spark和Flink批任务的生产集群里,Spark任务的资源“收放”更自如,Flink任务则经常出现“占着茅坑不拉屎”的情况。当然这不是说Flink一定不能跑批,只是说明它在这条路上还需要继续努力。
4. 生态、数据湖和周边工具的差距:批处理需要的不只是引擎
4.1 Hive语法兼容和SQL语义的细节落差
很多现有数仓底层是Hive,Spark能无缝地跑HiveQL,大部分Hive SQL可以直接迁移过来跑,这也是Spark这么多年来作为离线数仓核心引擎的基础。Flink做Hive兼容做了不少工作,支持Hive Catalog、Hive UDF等,但实际迁移中你会发现很多细节对不上。
举几个常见例子:
- Hive分桶表的写入和读取,Spark能很好兼容,Flink在某些版本下对分桶语义的支持并不完整。
- 自定义SerDe:大量老Hive表使用了自定义SerDe,Spark能直接读取,Flink在部分SerDe场景下会出现兼容性问题。
- SQL隐式转换和日期时间语义:Spark和Flink对同一SQL表达式的解析结果可能不一致,比如timestamp的精度、字符串和日期隐式转换的规则差异,这会导致同样的数据在两个引擎里算出来的结果不一样。
- Hive自定义UDF/UDAF:Spark的Hive UDF兼容性已经非常成熟,Flink也支持,但UDAF的迭代器调用、中间聚合态的序列化等场景就需要更多验证,我见过不止一次UDF在Flink里行为不一致引发的数据问题。
这类“细节落差”在批处理迁移时最容易踩坑,而且通常不会在一开始暴露,而是在某种特殊的数据组合下突然爆发,排查成本很高。
4.2 数据湖三件套:Spark是亲儿子,Flink更像是邻居家的孩子
现在做批处理基本绕不开数据湖技术:Hudi、Iceberg、Delta Lake。坦白讲,Spark在这三套数据湖体系中的地位是最核心的,几乎所有数据湖框架都优先支持Spark的读写和优化操作。
Spark在数据湖场景里的能力非常完整:小文件合并、Clustering、Z-Order、Incremental Read,这些数据湖维护操作在Spark侧都有成熟的命令和工具,批式读写性能也更好。Flink侧当然也提供数据湖的写入支持,尤其是实时写入场景,Flink的表现很好——但到了批读、批量更新、小文件治理这类批处理刚需时,Flink的能力就明显弱了一档。
一个典型的体验是:Iceberg在Spark上的批式读取可以自动应用文件裁剪、列裁剪和谓词下推,读取速度快且稳定;同样的Iceberg表在Flink上做批读,性能和稳定性往往需要额外调优,而且有些优化手段还不完全支持。这在湖仓一体架构里是非常要命的,因为湖仓一体本身就要求批读路径的强大支撑。
4.3 周边工具链:ML、图计算、调度框架适配
批处理从来不是只用SQL跑个数就完事,现代数据平台还需要机器学习、图计算、数据探索等能力。
Spark有MLlib、GraphX、Spark NLP等庞大的批处理生态,配合Notebook、Livy、Zeppelin等工具,可以做非常多的事情。Flink虽然也在演化Flink ML,但成熟度和覆盖范围跟Spark生态完全不在一个量级。做纯批处理的机器学习特征工程、模型训练任务,目前几乎没有人会用Flink替代Spark。
另外还有集群和调度平台的适配度问题。YARN、Kubernetes、自研调度平台对Spark批任务的接入经验都非常丰富,社区里的Spark集群搭建教程、参数优化资料一抓一大把。Flink批任务在调度平台的接入上就相对粗糙,很多平台对Flink批任务的支持仅仅是“能跑”,但监控、告警、资源隔离、任务优先级等方面都不如Spark顺手。
4.4 实际开发里绕不开的“小毛病”
这部分纯粹是实战中发现的一些琐碎但磨人的问题,也许每个单独看都不大,但堆在一起就能拉开体验差距:
- 读JSON文件:Spark的
inferSchema很好用,能自动推导嵌套结构;Flink SQL读JSON通常要手写复杂嵌套的ROW类型定义,开发效率直接打折。 - JDBC连接器:Flink的JDBC Sink在批任务里高并发写数据库时,容易出现连接超时、连接池不够用的情况,需要调
connectionPoolSize、socketTimeout等参数;Spark的JDBC写入路径经过多年打磨,行为的可预测性更强。 - 文件格式转换:把Parquet改成ORC,或者压缩算法换一下,Spark里DataFrame读写两行代码搞定;Flink里你要调整Sink的Table Store格式或者自己写FileSink,麻烦不少。批处理中经常要做的“改表格式、改文件布局”这种脏活累活,Spark显然更顺手。
- Flink UI看批任务不如Spark直观:Spark UI可以看到每个Stage的详细执行情况、Shuffle读写量、GC时间,排查问题非常高效。Flink UI在批任务执行计划可视化、阶段拆分上做得相对较弱,很多时候要看日志和反压监控,排查效率低一些。
5. 实战踩坑清单:Flink批处理中遇到的典型问题
5.1 同样的SQL,Flink比Spark慢了一半
这个坑我踩过不止一次。现象是:同样的数据量、同样的SQL,Spark十几分钟跑完,Flink要跑半小时甚至更久。排查下来通常是两个原因:一是Shuffle并行度过低或分区不均衡,二是SQL执行计划不够优。
解决思路也不复杂:检查Flink SQL执行计划,看Join策略、分区数、是否存在不合理的Broadcast或过多小文件。针对倾斜场景,手动做加盐和二次聚合;把table.exec.resource.default-parallelism调到一个合理的值,避免所有算子共用同一个默认并行度。但很遗憾,这些操作都需要DBA手工介入,不像Spark AQE那样自动解决。
5.2 Flink批任务OOM,内存开销比Spark大
Flink批任务OOM的坑,往往是内存模型差异导致的。Flink的TaskManager有托管内存(Managed Memory)、网络缓冲、堆外内存等区分,参数设置不合理就容易出现OOM或反压。
我的经验是:跑批时调大taskmanager.memory.managed.size,给排序、Hash聚合和Join留足空间;网络缓冲也要适当放大,否则大规模Shuffle时会成为瓶颈。相比之下,Spark的Executor内存模型经过这么多年的优化,文档和社区经验都极其丰富,调优路径更清晰。Flink的内存参数虽然文档也在完善,但很多细节还是得自己试。
5.3 从Spark迁移到Flink后,数据处理结果不一致
这是比较隐蔽的问题。数据结果不一致,通常是SQL语义差异引起的,比如时间函数、空值处理、字符串隐式转换规则不一致。例如Spark里date_add函数的取值范围、Flink对空字符串和NULL的区分方式,都可能和Hive不同。
这方面没有取巧的办法,只能一条一条SQL对照验证。我的建议是:在迁移前建立一套自动化对跑框架,同一输入数据分别在Spark和Flink上执行,输出结果做Diff,提前把语义差异暴露出来,千万别上线后再去排查数据质量问题。
5.4 快速定位问题的排查路径
如果Flink批任务出现性能问题,我个人建议按以下顺序排查:
- 先看执行计划,是否有多余的Shuffle、不合理的Join策略;
- 再看并行度设置,是否所有算子都使用同一并行度,倾斜源是否没有打散;
- 再看内存配置,TaskManager托管内存和网络缓冲是否够用;
- 最后看数据倾斜,重点检查大Key分布,必要时手动加盐。
这套排查思路,其实也从侧面说明了Flink批处理在自动化调优上的短板——太多本该平台做的事情,最后都落到了人头上。
6. 一张表格看清:批处理场景下Spark和Flink怎么选
6.1 选型判断速查表
| 场景 | 推荐引擎 | 原因 |
|---|---|---|
| 大规模离线ETL、复杂多表JOIN | Spark | 优化器成熟,AQE自动处理倾斜和分区合并 |
| 数据湖批式读写和湖表维护 | Spark | 数据湖框架对Spark支持最完整 |
| 机器学习、特征工程、图计算 | Spark | MLlib/GraphX生态无可替代 |
| 流批一体、实时数仓 | Flink | 流处理能力强,批侧能力在持续补齐 |
| Hive存量任务迁移 | Spark | Hive语法和UDF兼容性更成熟 |
| 实时写入数据湖 | Flink | 流式写入和状态管理优势明显 |
| 团队已有较强Flink技术栈 | Flink | 解决维护一套引擎的成本问题 |
| 运维监控工具链要求高 | Spark | 生态工具更丰富、排障经验更多 |
6.2 正在发生的变化:Flink并没有躺平
我也要说句公道话,Flink在批处理上的短板,社区是有认知的,而且正在努力追赶。Flink 1.18引入了新的批执行模式,优化器也在向CBO方向推进,自适应调度、Shuffle优化都在持续做。
Flink 2.0也在规划中,包括动态资源分配、更好的批式Shuffle、增强的SQL优化能力等。Paimon项目也在补数据湖批读批写方向的能力。整体趋势在变好,但往前看两三年,纯批处理场景下Spark的领先优势大概率还会保持。
6.3 我的建议:不要为了“统一”而强行替换
技术选型这事,最怕跟风。如果你现在的批处理主力是Spark,集群稳定、任务跑得好,就没必要因为“流批一体”的口号全线迁Flink。反过来,如果你们已经有了一套成熟的Flink流处理体系,批任务量不大,为了省运维成本,用Flink跑批也是可以接受的——前提是你做好踩坑的心理准备。
我个人在实践中的体会是:流批一体是长期方向,但现阶段“批归Spark、流归Flink”依然是最稳妥的架构。等Flink批处理优化器真正赶上Spark的那一天,我们再聊全面替换也不迟。现阶段,认清各自的能力边界,把合适的工作交给合适的引擎,才是性价比最高的做法。