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

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能否解析并优化执行逻辑:

  1. DSL表达式$"i" === 2的处理逻辑

    • $"i"是Spark提供的Column表达式,属于Spark SQL的DSL领域特定语言。这种写法直接操作Spark的逻辑计划节点,Catalyst优化器可以完全解析逻辑:
      • 识别出仅需列i即可完成过滤和后续select操作;
      • 执行列裁剪(Column Pruning),只读取存储中的i列;
      • 还能将过滤条件下推到数据源(PushedFilters),在读取阶段就完成过滤,减少IO开销。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 11:03:16