Spark DataFrame筛选/排除操作最优方案:两种方式对比
两种Spark DataFrame拆分方式的效率对比分析
结论很明确:第二种方式(先筛选小数据集,再用原表排除小数据集得到大数据集)的效率远高于第一种,你的直觉判断是错误的,核心原因在于except操作的成本逻辑:
except操作的本质是对两个数据集做全量去重比对,需要将两个数据集的所有数据进行哈希计算或排序后匹配,其计算、IO开销与参与比对的数据集大小直接正相关。
第一种方式的问题
你先筛选出大数据集large并持久化,再用dfMain.except(large)获取小数据集。这个过程中,except需要让整个原大表和另一个同量级的大表做全量比对,相当于处理2倍规模的大数据集,哈希计算、数据shuffle、磁盘IO的成本都会拉到最高,哪怕large已经持久化,也无法抵消这种量级的比对开销。
第二种方式的优势
先筛选出小数据集small并持久化,再用dfMain.except(small)获取大数据集。此时except仅需要让原大表和一个极小的数据集做比对:
- 小数据集
small的持久化成本极低(内存/磁盘占用小); - 比对过程中只需要加载小数据集的哈希值或全量数据,和大表做快速过滤,shuffle和计算量都被降到最低。
额外优化建议
其实还有更高效的拆分方式——直接用两次filter分别获取两个数据集,完全避免except的开销:
val large = dfMain.filter(col("data") === "criteriaA") val small = dfMain.filter(col("data") === "criteriaB") large.write.parquet("storeA") small.write.parquet("storeB")
如果criteriaA和criteriaB是互斥且覆盖所有数据的场景,还可以通过一次数据扫描生成两个数据集(比如利用Spark的侧输出流),进一步减少IO次数。
内容的提问来源于stack exchange,提问作者Jelly
相关产品推荐
相关产品推荐

