Spark不同过滤方法的性能差异及实现方案对比咨询
问题描述
我有一个大型DataFrame(记为rows_df),需要为其应用多个过滤条件(filter1、filter2、filter3、filter4)。当前实现方法如下:
filtered_df1 = rows_df.filter(condition1) filtered_df2 = rows_df.filter(condition2) filtered_df3 = rows_df.filter(condition3) filtered_df4 = rows_df.filter(condition4) filtered_df1_wc = filtered_df1.withColumn(filter1, "1") filtered_df2_wc = filtered_df2.withColumn(filter2, "1") filtered_df3_wc = filtered_df3.withColumn(filter3, "1") filtered_df4_wc = filtered_df4.withColumn(filter4, "1") filtered_df = filtered_df1.union(filtered_df2).union(filtered_df2).union(filtered_df3).union(filtered_df4) filtered_df_with_col = filtered_df1_wc.union(filtered_df2_wc).union(filtered_df2_wc).union(filtered_df3_wc).union(filtered_df4_wc) good_rows = rows_df.exceptAll(filtered_df) bad_rows = filtered_df_with_col
三种替代实现方案如下:
Alternative 1:
for each filter: rows_df = rows_df.withColumn(filter1, when(condition1,1).otherwise(0))
Alternative 2:
rows_df = rows_df.withColumn(filter1, when(condition1,1).otherwise(0)) .withColumn(filter2, when(condition2,1).otherwise(0)) .withColumn(filter3, when(condition3,1).otherwise(0)) .withColumn(filter4, when(condition4,1).otherwise(0))
Alternative 3:
select *, case when condition1 then 1 else 0 end as filter1, case when condition2 then 1 else 0 end as filter2, case when condition3 then 1 else 0 end as filter3, case when condition4 then 1 else 0 end as filter4 from rows_df
所有替代方案均可通过生成的标记值区分有效行(good_rows)与无效行(bad_rows)。我认为当前方法性能较差,因需依次执行过滤,Spark要扫描DataFrame4次。请问替代方案是否也存在该问题?哪种方案性能更优?另外,Alternative1与Alternative2是否存在差异?我认为前者每次迭代生成新DataFrame,后者一次性执行所有条件,是否正确?
解答
1. 替代方案是否存在重复扫描问题?
当前方法确实会触发4次全表扫描,因为每次filter操作都会单独读取原DataFrame的数据。而三种替代方案都只会触发1次全表扫描:
不管是链式调用withColumn、循环调用withColumn,还是Spark SQL的单条查询,Spark的Catalyst优化器都会将这些操作合并为一个执行计划,只需要扫描一次原始数据,在扫描过程中同时计算所有过滤条件的标记值。
2. 哪种方案性能更优?
三种替代方案的性能差异极小,因为最终生成的执行计划几乎一致:
- Alternative 2和Alternative 3本质等价,都是一次性声明所有列的计算逻辑,优化后的执行计划完全相同。
- Alternative 1虽然是循环调用
withColumn,但Spark的惰性求值机制会将这些连续的withColumn操作合并,不会额外增加扫描次数,最终执行计划和Alternative 2一致。
细微差距层面:
- Alternative 2的DataFrame链式调用写法,在代码可读性和维护性上更优,无需切换到SQL语法。
- 如果过滤条件本身就是SQL表达式,Alternative 3使用起来更方便,但需要先将DataFrame注册为临时视图。
3. Alternative1与Alternative2是否存在差异?
你的理解不完全准确:
- 两者都会生成多个DataFrame对象,但Spark是惰性求值的,这些DataFrame只是执行计划的逻辑节点,不会立即触发计算。
- 循环调用
withColumn(Alternative1)和链式调用(Alternative2)最终生成的执行计划完全相同,Catalyst会把所有withColumn操作合并成一个阶段,一次性完成所有标记列的计算,不会因为循环产生额外开销。 - 唯一区别是代码写法:Alternative2更紧凑,Alternative1适合过滤条件数量不确定(比如从配置动态读取条件)的场景。
内容的提问来源于stack exchange,提问作者Vandit Goel
相关产品推荐
相关产品推荐

