Spark使用Scala的_.i语法过滤时为何会读取全列?
问题背景
以下代码表现良好,仅读取列i(注意最后一行ReadSchema: struct<i:bigint>):
import org.apache.spark.sql.Dataset // Define the case class case class Foo(i: Long, j: String) // Create a Dataset of Foo val ds: Dataset[Foo] = spark.createDataset(Seq( Foo(1, "Q"), Foo(10, "W"), Foo(100, "E") )) // Filter and cast the column val result = ds.filter($"i" === 2).select($"i") // Explain the query plan result.explain() // It prints: //== Physical Plan == //*(1) Filter (isnotnull(i#225L) AND (i#225L = 2)) //+- *(1) ColumnarToRow // +- FileScan parquet [i#225L] Batched: true, DataFilters: [isnotnull(i#225L), (i#225L = 2)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[dbfs:/tmp/foo], PartitionFilters: [], PushedFilters: [IsNotNull(i), EqualTo(i,2)], ReadSchema: struct<i:bigint>
但如果使用val result = ds.filter(_.i == 10).map(_.i),物理执行计划会读取包括j在内的所有列(注意最后一行ReadSchema: struct<i:bigint,j:string>):
//= Physical Plan == //*(1) SerializeFromObject [input[0, bigint, false] AS value#336L] //+- *(1) MapElements //$line64a700cfcea442ea899a5731e37978a9115.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$Lambda$8811/2079839768@1028cff, obj#335: bigint // +- *(1) Filter //$line64a700cfcea442ea899a5731e37978a9115.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$Lambda$8810/701415521@212ee011.apply // +- *(1) DeserializeToObject newInstance(class //$line64a700cfcea442ea899a5731e37978a925.$read$$iw$$iw$$iw$$iw$$iw$$iw$Foo), obj#334: //$line64a700cfcea442ea899a5731e37978a925.$read$$iw$$iw$$iw$$iw$$iw$$iw$Foo // +- *(1) ColumnarToRow // +- FileScan parquet [i#225L,j#226] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[dbfs:/tmp/foo], PartitionFilters: [], PushedFilters: [], ReadSchema: struct<i:bigint,j:string>
请问为何Spark在filter中使用Scala的_.i语法时,处理方式会有所不同?
原因解析
两种写法的核心差异在于Spark能否解析并优化执行逻辑:
DSL表达式
$"i" === 2的处理逻辑$"i"是Spark提供的Column表达式,属于Spark SQL的DSL领域特定语言。这种写法直接操作Spark的逻辑计划节点,Catalyst优化器可以完全解析逻辑:- 识别出仅需列
i即可完成过滤和后续select操作; - 执行列裁剪(Column Pruning),只读取存储中的
i列; - 还能将过滤条件下推到数据源(PushedFilters),在读取阶段就完成过滤,减少IO开销。
- 识别出仅需列
Scala匿名函数
_.i == 10的处理逻辑_.i是Scala语法糖,本质是生成匿名Lambda函数,输入为Foo类型的JVM对象。Spark无法解析Lambda内部逻辑,只能按要求先实例化完整的Foo对象才能执行操作:- 为构造
Foo对象,必须读取所有列(包括j),因为Foo类的构造依赖i和j两个字段; - 执行流程变为:读取全列 → 反序列化成
Foo对象 → 调用Lambda过滤 → 提取i字段序列化输出; - 这种场景下,Catalyst优化器无法做列裁剪和过滤下推,只能被动执行对象层面的操作。
- 为构造
简单来说,DSL写法是明确告诉Spark操作的列和逻辑,让它针对性优化;Lambda写法则是要求Spark把完整对象交给用户处理,只能加载全量数据。
内容的提问来源于stack exchange,提问作者jack
相关产品推荐
相关产品推荐

