Flink批处理 vs Spark:六大短板与选型建议 最近我在好几个技术社群里连续被问到同一个问题团队想推流批一体准备把批处理也迁到 Flink 上到底值不值说 Flink 流处理强大家基本没意见但一聊到批处理争论就多起来了。真去翻社区的吐槽最常看到的就是“Flink 批处理慢”“Flink 一跑批就 OOM”“还是 Spark 香”。这些说法有道理但也不全对。这篇文章我就从工程实践的角度把“在批处理方面Flink 相对于 Spark 还有哪些不足”这个问题系统拆一遍帮助正在做流批一体选型的技术负责人、准备把批任务迁到 Flink 的工程师以及想在两个引擎之间做判断的入门同学都能搞清楚差距到底在哪、哪些是可以接受的、哪些是真坑。1. 先给这场对比定个坐标Flink 的批处理到底做成什么样了1.1 技术基因决定起点一个流出身一个批出身讨论 Flink 批处理不足先得明白 Flink 是怎么一路走过来的。Flink 最初是为流处理设计的早期版本里的 DataSet API 是独立的批处理 API和 DataStream API 共享执行引擎但优化器、执行计划、API 语义完全是两套东西。那个阶段谈批处理大家只会想到 Spark、Hive on TezFlink 的批处理在社区里基本没有存在感。2020 年前后Flink 开始喊“流批一体”核心思路是把批处理当作有界流的特例来统一。1.12 版本里 DataStream API 支持了 BATCH 执行模式Table API 和 SQL 成为批处理的主推入口DataSet API 进入弃用周期。1.14、1.15 又把 sort-shuffle、blocking shuffle 等批处理优化陆续引入。但这一整套补课从时间线上比 Spark 晚了好几年而 Spark 从 RDD 到 DataFrame 再到 SQL批处理相关的优化器、shuffle、内存管理已经迭代了非常多个版本。这个“出身差异”带来一个很实际的影响Flink 的批处理能力和 Spark 不在同一个成熟周期上。不是说 Flink 做不了批而是它做到什么程度、踩坑之后的文档和社区经验积累都比 Spark 少一圈。团队里如果没人深度用过 Flink 批处理第一年往往会比较难受。1.2 Flink 为了“能跑批”做了哪些补课这里简单过一下 Flink 在批处理上的几个关键节点方便后面理解哪些短板是新补的、哪些还有坑。1.12 版本引入 BATCH 执行模式让 DataStream 可以按有界流来跑批从底层统一了流批。但注意早期 BATCH 模式下的很多算子行为跟真正的批引擎还不一样比如 SQL join、sort 的实现是在后面几个版本里才逐步替换掉的。1.13 到 1.15 版本Flink 把 shuffle 机制从普通网络缓冲改成 sort-shuffle也就是先对数据做排序再落盘类似 Spark 的 sort-based shuffle减少小文件问题对大规模 shuffle 的稳定性提升明显。同时 blocking shuffle 在 BATCH 模式下取代流式的 pipelined shuffle解决了批任务里下游还没执行完、上游数据就已经推过去的问题让调度和资源利用更接近批引擎。2021 年前后DataSet API 被标记为废弃Table/SQL 成为 Flink 批处理的唯一正门。这带来一个好处就是开发体验向 Spark SQL 看齐但也带来一个难受的点如果需要算子级别的细粒度操作比如自定义分区、精确控制 shuffle 方式在 Table API 里会别扭很多很多在 Spark 里可以用 RDD/DataFrame 灵活做的事情Flink 里没有那么多的选择。这些补课让 Flink 批处理能跑、能用到生产但“能跑”和“在同等规模复杂度下比 Spark 更稳更快”是两个完全不同的概念。1.3 衡量批处理能力的五个维度我平时做引擎对比一般不看社区里的玄学评价只看五个维度稳定性也就是会不会 OOM、会不会跑挂吞吐与性能同样资源谁跑得快SQL 能力与易用性复杂查询能不能写、优化器聪明不聪明生态连接器需要的源和目标是否开箱即用运维复杂度排障、调优、扩缩容是否顺手。后面拆解 Flink 的短板基本就按这条线走。2. 逐个拆解Flink 批处理相对 Spark 的短板在哪这一章是全文的核心我分六个小节讲每个点都是我自己在项目里真实碰到过的不是从文档里抄来的。2.1 内存模型堆外内存就是一个大坑Spark 的内存管理虽然也被吐槽过executor 内存参数确实不好调但它的 off-heap 是可选配置默认大部分算子数据走堆内存shuffle 数据走磁盘加堆内缓冲整个模型相对直观。Flink 完全不同TaskManager 的内存体系非常细堆内、堆外、托管内存、网络缓冲、JVM Metaspace、Direct Memory 全部是分开配置的。我第一次用 Flink 跑一个中等规模的批任务时就踩了典型的堆外内存超限问题。任务本身不算大几十 GB 的输入数据做了几个 join 和 group by结果跑了一个多小时之后 TaskManager 直接挂掉日志里报的是 Direct buffer memory 超限。排查了一轮才发现Flink 的托管内存默认会占进程内存的一半左右而排序、hash join、RocksDB 都要用它。批处理里的排序聚合比流处理重得多如果没主动调 taskmanager.memory.managed.size很容易在排序、join 的中间阶段把堆外内存打爆。从 Spark 迁到 Flink 批处理第一个要适应的就是不要再去调 JVM 的 -Xmx而是要理解 taskmanager.memory.process.size、taskmanager.memory.jvm-metaspace.size、taskmanager.memory.managed.fraction 这一整套体系。这不算 Flink 的硬伤但对从 Spark 过来的团队来说学习成本和排障成本实实在在比 Spark 高一截。你可以抱怨“又是堆外内存”但更重要的还是掌握解法一是调大 managed memory 并限制 JVM overhead二是给排序和 join 算子设置合理的内存上限三是确认 RocksDB 的 block cache 配置没有和 managed memory 打架。2.2 调度模型单作业调度吃大亏Spark 的 DAG Scheduler 把一个作业拆成多个 stagestage 内部又能并行调度成百上千个 task。更重要的是Spark 是“一个 Application 多个 Job”的模式在一个 app 里可以连续跑多个动作每个 job 之间共享 executor 资源。Flink 的调度模型本质上还是“一个作业一张 ExecutionGraph”即使在 BATCH 模式下也是绑定在一个 JobManager 上。假如你有 100 个小批任务Spark 可以塞到一个 application 里依次执行复用 executorFlink 则要么写成一个作业的不同 branch要么起 100 个作业每个作业各自申请资源、各自调度JobManager 的压力也大。实际项目里用 Spark 做大量小 ETL 任务、报表任务非常自然而 Flink 要吃到同样的吞吐要么靠共享 slot 的配置去凑要么就得接受任务粒度的资源浪费。还有一个细节Spark 的动态资源分配可以根据 stage 执行情况伸缩 executor闲置时归还资源Flink 虽然也支持一些弹性伸缩能力但主要面向流任务批处理场景下我很少见到有人真把它用到生产。这导致在“一堆大小不一、混合提交”的批任务场景里Flink 对集群资源的利用率通常不如 Spark。当然这不是说 Flink 调度一无是处它延迟低、对流任务友好但对批任务的批量化调度模型确实没那么匹配。2.3 Shuffle 与网络传输批场景的吞吐天花板这是很多人在讨论 Flink 批处理性能时容易忽略的关键点。Spark 的 shuffle 是强落盘加 sort 的shuffle 数据写到本地磁盘fetch 端再从对面拉配合 shuffle service 的文件索引在几个 TB 级别 shuffle 的任务里表现很稳定。Flink 早期是流式 shuffle数据在内存和网络缓冲区里直接推给下游这在流任务实时性上是优势但批任务里如果下游还没起来、或者数据积压就会很难受。Flink 在 BATCH 模式下改用了 blocking shuffle把上游输出先落盘下游再读取思路跟 Spark 接近了。但落盘和读取的实现细节还是有差距Spark 的 shuffle 文件管理、索引、聚合 fetch 做了多年优化Flink 的 sort-shuffle 是在 1.15 之后才逐步稳定的。如果你打开 Flink 批任务的 Shuffle Service 监控会发现小文件多、读放大的时候性能和稳定性都要打折。我自己在一个 10 节点测试集群上做过一个非严格对比两个引擎处理大约 200GB 的排序聚合Spark 用 sort-based shuffleFlink 开 BATCH 模式加 sort-shuffle在同样并行度下Flink 的 shuffle 耗时大概多出 20% 到 40%。这个数字不构成严谨基准但能反映一个趋势在最吃 shuffle 的批场景里Flink 的优化底子还在追赶。2.4 数据倾斜与动态优化没有 AQE 很被动Spark 3.x 引入的 AQE 真的是批处理里的一个重要改进。它可以动态合并 shuffle 分区、动态切换 join 策略、自动处理 join 中的数据倾斜。做过数据开发的人应该都知道 Spark 的 skew join 自动优化一旦检测到某个 key 数据量特别大AQE 会在 shuffle 之后重新聚合热点 key自动加盐拆分。Flink 这边即使到了近期的版本也还没有一个跟 AQE 对等的、能自动处理倾斜的机制。官方建议的倾斜处理方式基本是靠 SQL hint、自定义分区器或者手动给 key 加随机后缀全部是人工干预。做一次两次还能忍但如果你的数据分布是每天变的、倾斜位置不固定的批任务调优就变成了一场持久战。我之前做一份用户画像离线计算任务维度表很大分组之后的某些 key 数据量非常大Flink 批模式跑起来之后个别子任务迟迟不结束整个作业就等它。排查要通过 Web UI 或者 Metrics 定位是哪个算子的哪个 subtask 慢再针对性地改并行度或加盐。换成 Spark 开 AQE同样的场景大概率自己就处理了你要做的主要是上线前确认和效果观察。这个差距两边的用户体会会很深。2.5 SQL 与优化器功能追赶中差距仍可感知Flink SQL 这些年进步确实大1.16、1.17 之后 planner 基本重写过支持了更多优化规则比如 join 顺序优化、谓词下推、投影裁剪。但和 Spark SQL 的 Catalyst 优化器相比规则丰富度和成熟度上还有可感知的差距。我在实际项目里碰到过的几个场景复杂多表 join 的自动顺序调整。Spark 靠 CBO 基于统计信息调 join 顺序Flink 虽然也有代价优化但需要你主动提供更多统计信息很多时候优化结果还是需要手动重排 SQL 或者加 hints 引导。子查询和窗口函数。一些很复杂的嵌套子查询、长时间范围的窗口计算Spark SQL 写起来比较顺手Flink SQL 偶尔会因为语法限制或者优化不到位需要拆成多段 SQL 加临时表处理。动态分区插入和分区裁剪。Spark 对动态分区处理、静态分区裁剪更成熟。Flink 在写 Hive 表的分区处理上过去没少出问题最近几个版本才稳定下来。还有一个很实际的痛点Flink SQL 对某些批处理场景里常用的写法兼容性不如 Spark。如果你的团队里写 SQL 的人大多是 Spark 或者 Hive 出身迁到 Flink 之后写简单查询问题不大一写复杂查询就会感觉处处受限。2.6 连接器与生态连接器不少坑也不少Flink 的连接器数量并不少官方维护的 JDBC、Kafka、Hive、Hadoop FileSystem、Iceberg、Hudi、Paimon 都有。不过论“连接器在批处理场景下的成熟度”Flink 离 Spark 还有距离。一个典型例子是 JDBC 连接器这也是最近各个群里经常被吐槽的内容Flink 的 JDBC 连接器在批模式下批量写 MySQL 或 PostgreSQL默认的 sink.buffer-flush.interval 是 1 秒如果不调大 batch size写入延迟会很高但 batch size 一旦设得太大事务缓存超过阈值又会报错。再加上 upsert 支持不够顺滑要实现“有则更新、无则插入”要么绕道用 executor 模式要么自己写 sink。Spark 的 JDBC 连接器虽然也不算完美但写入的并发控制、重试机制、批量提交在大多数版本里表现都稳定不少。再说一个常见场景读 JSON 文件。Spark 用 spark.read.json 一行搞定支持多行 JSON、schema 推断、嵌套结构Flink SQL 读 JSON 文件建表时就得声明 json format 的 schema遇到单行 JSON 的复杂嵌套有时还要额外写一堆 JSON 函数去处理。你可以搜一下“spark 中读取 json”的经验贴几乎全是 Spark 的做法这侧面说明在文件类数据源处理上社区的默认经验值还是 Spark 更顺手。Hive 集成也一样。Flink 对 Hive 的读写支持有1.15 之后也能走 HiveCatalog但一些细节比如分区发现、文件格式兼容、Hive UDF 调用、复杂类型支持实际使用中还是会踩到。我们之前用它读 Hive 表遇到过 Hive 分区读取时类型推断不一致的问题排查了大半天。Spark 这边由于和 Hive 同源兼容性天然好很多。当然这也分场景不是所有连接器都是 Flink 差比如 Kafka、Paimon 这种 Flink 深度参与的项目反而是 Flink 的体验更好。3. 用对比实验看差距一次真实的批任务选型记录3.1 实验同一份 SQLSpark 和 Flink 各自的配置与表现光说概念容易飘我分享一个可复现的对比实验思路。数据建议用 TPC-H 生成规模选 50GB 到 200GB 之间比较合适集群配置按 4 个 worker 节点、每个节点 8C16G 来搭。如果你还没有现成环境可以参照网上的 Spark 集群搭建教程先把基础环境跑起来但引擎层面的配置差异很大要分别处理。具体步骤大致是这样用 tpch-dbgen 生成 TPC-H 数据scale factor 可以取 30 或者 100对应大概 30GB 或 100GB 原始数据。Spark 侧把同样的 TPC-H SQL 在 spark-sql 里执行开启 AQEexecutor 内存给 4G每个 executor 2 核动态资源分配关掉保证资源固定。Flink 侧execution.runtime-mode 设为 BATCH并行度设置与 Spark 总 executor 核数接近taskmanager 内存按照前面的建议调好SQL 用 Table API 的 statement set 提交。统一记录查询耗时、shuffle 数据量、任务稳定性看有没有发生 OOM 或重试。我做过一次非严格对比结果大致可以用下面的表表达注意这是趋势性观察不是严谨基准引擎查询类型数据规模耗时表现稳定性表现Spark单表聚合50GB基准稳定Flink BATCH单表聚合50GB与 Spark 相近稳定Spark多表 join 大 shuffle200GB基准稳定AQE 自动调优Flink BATCH多表 join 大 shuffle200GB比 Spark 慢约 20%-40%需要手动调内存否则有 OOM 风险这个实验最有价值的结论是Flink 不是所有批任务都慢。单表扫描、简单过滤聚合Flink 和 Spark 基本打平一旦涉及大 shuffle、多 join、倾斜数据Flink 的调优成本和性能差距就显现出来了。做选型的时候不要笼统说“Flink 批处理不行”要拿自己的业务数据去跑一遍用这类实验说话。3.2 Flink 批任务的调参清单与执行模式如果你决定试一把 Flink 批处理下面这份调参清单可以直接抄作业都是我在生产环境里验证过的基础配置execution.runtime-mode: batch强制走批执行计划。taskmanager.memory.process.size按容器内存设比如 8G 就写 8g。taskmanager.memory.managed.fraction建议 0.5 到 0.6批处理排序、join 对托管内存需求大。parallelism.default按可用 slot 数来别超过 total cores 太多。taskmanager.numberOfTaskSlots建议 2 到 4不要每个 slot 都塞一个吃满内存的任务。table.optimizer.join-reorder-enabled: true打开 join 重排。如果任务里有大状态state.backend 用 rocksdb并单独控制 block cache 大小。这里特别要说明为什么 slot 不要设太多。批任务里一个算子一个并行实例就是一份执行资源slot 设得多并行度自然高但如果每个任务都吃内存托管内存会迅速被瓜分干净后续排序、shuffle 就频繁溢写磁盘性能直线下降。Spark 的做法是让一个 executor 里跑多个 task 共享内存Flink 的 slot 隔离粒度更细配置思路要反过来并行度够用就行内存要给足。还有一个执行模式上的坑批任务的 checkpoint 默认不会自动打开因为不需要流式状态恢复。但如果一个跑了两小时的大批任务中间失败没有 checkpoint 就得从头来过。想让批任务有断点续跑能力要显式开启 checkpoint间隔可以设大一些比如 5 到 10 分钟这样成本可控又能拿到恢复能力。3.3 从 Spark 迁移到 Flink 批处理你一定会碰到的差异这里把我在迁移过程中遇到的高频差异整理成对照表能帮你快速对齐能力点SparkFlink核心 APIRDD / DataFrame / SQLDataStream / Table API / SQL执行单位Application 内多 Jobstage 调度一个作业一张 ExecutionGraph内存模型executor 堆内为主off-heap 可选TaskManager 堆外托管内存为主配置复杂Shufflepush 式 sort-based shuffleBATCH 模式 blocking shuffle sort-shuffle自适应优化AQE可自动处理倾斜、动态调整 join无对应机制需手动调优连接器成熟度文件、JDBC、Hive 生态成熟稳定核心连接器可用但部分场景细节有坑SQL 优化器Catalyst CBO规则丰富Planner 持续迭代复杂场景仍需手动优化从开发习惯上最直接的差异有两个。第一Spark 里常见的 mapPartitions、broadcast 变量在 Flink Table API 里没有直接等价物要么拆到 DataStream API 的 process function 里实现要么用 SQL hint 来做广播 join。第二Spark 读文件自动推断 schemaFlink 建表时要把格式、字段、主键都写清楚写起来啰嗦但对生产环境来说更可控。还有一点容易被忽略空值和 null key 的处理。Spark 的 join 默认会过滤掉 null keyFlink 的行为在某些版本里不太一样如果关联键可空迁移后要显式确认 join 结果是否符合预期不然会出现结果行数不一致的问题。4. Flink 跑批的常见问题与排查技巧实录4.1 批任务 OOM别只会怪 Flink 吃内存先看堆外内存配置Flink 批任务跑挂十次有八次是堆外内存问题。现象通常是任务跑了一个多小时某个 TaskManager 突然死亡日志里出现 OutOfMemoryError可能是 Java heap space也可以是 Direct buffer memory。很多人一见 OOM 就以为是并行度开太高或者数据量太大其实多半是托管内存和 JVM overhead 的配置没有跟着批任务的真实需要走。排查顺序我建议这样去 TaskManager 日志里定位 OOM 类型是 Java heap space、Direct buffer memory还是 Metaspace。如果是 Direct buffer memory优先检查 taskmanager.memory.managed.size 和 jvm-overhead看是不是排序、shuffle 阶段堆外内存超限。如果是 Java heap space查算子的代码是不是在 DataStream 里攒了大量对象或者 join 时生成了过大的缓存对象。记录下 OOM 出现时的并行度、内存配置和任务阶段对比正常配置找差异。排查完之后别急着无限加内存先把并行度降下来试试。批任务的很多 OOM 是因为并行度太高每个 slot 分到的托管内存太少排序时全在溢写磁盘最后搞到 spill 和内存竞争一起爆发。4.2 BATCH 模式下 Sink 不写数或写不完是挺常见的事另一个常见坑是Flink 批任务明明跑完了下游表里却没有数据或者数据比预期少。很多人第一反应是 SQL 写错了实际原因往往出在 sink 的提交机制上。Flink 的很多 sink 是为流处理设计的基于 checkpoint 或两阶段提交来保证一致性。批任务如果不开启 checkpointsink 可能只在最后 flush 一次如果任务在 flush 前异常结束数据就丢了。如果你在写文件类 sink还会遇到 part 文件没有 merge 的情况看起来就是数据没写完整。我的处理建议是批任务如果对准确性要求高也把 checkpoint 打开间隔可以放宽到几分钟成本可控写文件时用 file system connector 的提交机制确认作业结束后触发文件提交和目录整理。排查时先看作业整体是不是 all succeeded再看 sink 算子的 metrics 里的写出条数能快速定位是哪一层丢的。4.3 JDBC 连接器在批任务里又慢又不稳原因和解决办法结合热词里被搜爆的“flink 的 jdbc 连接器异常”这里专门说一下。Flink 用 JDBC sink 写关系型数据库和 Spark 的体验差在默认参数上。默认的 sink.buffer-flush.interval 是 1 秒也就是说每秒才 flush 一次批量写入如果 batch size 又没调到足够大写入吞吐就特别难看。比较稳的配置是sink.buffer-flush.max-rows 设到 1000 或以上sink.buffer-flush.interval 设到 500ms 到 1s 之间并行度不要太高不然一个 checkpoint 周期内所有并发连接同时提交数据库很快就打满了。另外JDBC 连接器不支持真正的 upsert要实现主键冲突时更新需要目标表有主键并通过 SQL 自己拼 INSERT ... ON DUPLICATE KEY UPDATE这一点和 Spark 里的 JDBC 写入要额外处理是类似的但 Flink 的 jdbc connector 对事务重用的限制更明显。此外批任务里多个并行实例写同一张表连接数会线性增长。我们以前有个任务并行度 32直接把 MySQL 连接数打爆了。后来把并行度降到 8batch size 调大写入耗时反而降了一半。这里想提醒的就是偶发连接异常不一定是代码问题先看并发模型。4.4 倾斜一眼就能看出来但没有一键开关Flink 批任务的倾斜问题比 Spark 难处理得多。判断方法不算复杂看 Web UI 里每个 subtask 的 recordsIn 数量或者看反压状态某个 subtask 明显高于其他几个基本就是倾斜了。处理办法无非三板斧。第一把参与 join 的小表做 broadcast通过 SQL hint 强制走 broadcast join能避免大表 shuffle。第二对倾斜 key 做加盐拆分先在 SQL 里对 key 做 concat 随机后缀聚合完再按真实 key 汇总。第三如果倾斜发生在自定义算子用 DataStream API 定制分区器把热点 key 单独路由到多个 subtask。这三招在 Spark 里 AQE 已经能自动做大部分了在 Flink 里每次都要人肉处理。如果你手上有每天数据分布都在变的批任务我的建议是提前在同步或者预处理环节把热点 key 打散不要在引擎层临时救火。4.5 从零开始排障时Spark 优化经验迁移过来容易失效的两个场景最后提醒两个容易踩的迁移坑。第一个是参数迁移在 Spark 里调 spark.sql.shuffle.partitions 能显著影响性能Flink 里没有对应的全局 shuffle 分区参数需要调的是整体并行度或者对特定算子单独调 parallelism。如果照搬 Spark 思路去调 Flink基本没用。第二个是缓存与广播的迁移Spark 里缓存 DataFrame、广播小表非常顺手Flink 里 DataStream 的广播机制更重Table API 里的 broadcast join 又依赖 hint实际效果和表大小、集群配置关联很大。我在项目里试过强行对所有小表广播结果网络开销比 shuffle 还大整个任务更慢。正确做法是先检查小表的大小只有在明显小于广播阈值时才用 hint否则老老实实走正常 join。下面这张速查表方便你排查问题时快速定位问题现象大概率原因建议操作TaskManager 直接挂掉日志 Direct buffer OOM托管内存配置不足调大 managed.size降并行度检查 RocksDB cache批任务跑完目标表没数据sink 在批模式下未正确提交开启 checkpoint确认 sink connector 支持批提交JDBC 写入极慢或连接异常batch size 太小、并行度太高调大 flush-max-rows降并行度控制总连接数单个 subtask 明显拖慢整体任务数据倾斜使用 broadcast hint、加盐、自定义分区器手动处理SQL 同样的写法Flink 比 Spark 慢很多没有 AQE没有自动优化手动重排 SQL补充统计信息必要时换 join 类型排障经验说白了就是先看内存配置再看 sink 提交和并发模型最后才去怀疑 SQL 和算子这个顺序能省下很多时间。我个人在实际项目里的体会是Flink 批处理并不是不能用的替代品它最适合的场景是“你已经用 Flink 做核心流计算批数据规模在百 GB 级别、以 T1 跑批为主想省掉一套 Spark 引擎”。这种情况下Flink 的管理成本和数据链路统一带来的收益完全可以覆盖它在批处理性能上的落后。但如果你的批处理是核心业务量大、复杂 SQL 多、还要支撑分析师自助查询现阶段老老实实用 Spark 会更稳。最后分享一个小技巧不管最终选哪个引擎先把业务里最重的三张表和最复杂的三个查询拿出来单独跑一遍对比这个动作比看一百篇文章都有用。