为何减少分区数可避免Spark序列化结果过大引发的StageFailure?
orderBy(rand())触发序列化结果过大错误,减少分区后反而解决? 在Databricks(Spark 3.3.0 + Scala 2.12)环境中,从S3存储桶读取Parquet文件得到一个拥有662个分区的大型DataFrame,调用.orderBy(rand())生成随机排序版本时,频繁触发如下SparkException:
SparkException: Job aborted due to stage failure: Total size of serialized results of 36 tasks (1038.9 MiB) is bigger than spark.driver.maxResultSize 1024.0 MiB.
但先执行repartition(10)减少分区后,再调用.orderBy(rand()),错误就不再出现。按常理减少分区会增大每个网络传输块的大小,总序列化结果尺寸应该更大才对,想知道背后的原因。
代码示例
// 触发错误的代码 display(my_df .orderBy(rand()))
// 正常运行的代码 display(my_df .repartition(10) .orderBy(rand()))
这个现象的核心原因在于Spark任务结果序列化的额外开销,以及orderBy(rand())的执行逻辑差异:
序列化的分区固定开销
Spark把Executor的任务结果序列化返回给Driver时,每个分区的结果都会附带元数据(比如分区信息、序列化框架头部数据等)固定开销。原DataFrame的662个小分区会产生大量这类固定开销,累加后再加上数据本身的大小,总序列化结果就超过了spark.driver.maxResultSize的1024MB限制。
而repartition(10)把小分区合并成10个大分区后,固定开销的总数量从662份降到10份,即使每个大分区的数据量更大,但总序列化大小(数据+固定开销)反而控制在了限制以内。orderBy(rand())的执行流程差异
- 直接调用
orderBy(rand())时,Spark会先对662个分区分别做本地随机排序,再通过shuffle把所有分区的数据汇聚到一个节点完成全局排序,最后该节点需要把排序后的所有分区数据返回给Driver——662份带固定开销的序列化结果累加后触发了大小限制。 - 先执行
repartition(10)时,第一步已经把数据合并成10个大分区,后续orderBy(rand())只需要处理10个分区的本地排序和全局汇聚,最终返回给Driver的是10份序列化结果,总开销大幅降低,自然不会触发限制。
补充:原DataFrame的662个分区来自S3的Parquet小文件,小分区本身就会带来额外的序列化、调度开销,合并分区后不仅解决了本次错误,还能提升整体任务的执行效率。
内容的提问来源于stack exchange,提问作者Felipe

