You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 05:50:13