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

Spark中如何检测输入DataFrame里被转换逻辑实际使用的列?

检测Spark DataFrame中被实际使用的输入列

方法1:直接解析分析后的逻辑计划(推荐)

不需要触发Action,也不用解析文本格式的执行计划,直接通过操作Spark内部的逻辑计划结构提取被引用的列,可靠性和效率更高:

  1. 获取输出DataFrame经过分析后的逻辑计划
  2. 递归遍历计划节点,收集所有被引用的Attribute
  3. 关联这些属性到输入DataFrame的列,过滤得到实际被使用的列

示例Scala代码:

import org.apache.spark.sql.catalyst.plans.logical._
import org.apache.spark.sql.catalyst.expressions.Attribute

def collectUsedAttributes(plan: LogicalPlan): Set[Attribute] = plan match {
  case attrRef: AttributeReference => Set(attrRef)
  case unaryNode: UnaryNode => collectUsedAttributes(unaryNode.child)
  case binaryNode: BinaryNode => collectUsedAttributes(binaryNode.left) ++ collectUsedAttributes(binaryNode.right)
  case otherNode => otherNode.children.flatMap(collectUsedAttributes).toSet
}

// 假设inputDf是你的输入DataFrame,outputDf是最终转换后的输出DataFrame
val usedAttributes = collectUsedAttributes(outputDf.queryExecution.analyzed)
val usedColumns = inputDf.schema.fields
  .filter(field => usedAttributes.exists(_.name == field.name))
  .map(_.name)

println("实际使用的输入列:" + usedColumns.mkString(", "))

方法2:利用Spark的列裁剪优化结果

Spark内部的列裁剪(Column Pruning)优化会自动移除未被使用的列,直接查看优化后的逻辑计划中输入节点保留的列即可:

import org.apache.spark.sql.catalyst.plans.logical.LogicalRelation

def findUsedInputColumns(plan: LogicalPlan): Set[String] = plan match {
  case logicalRel: LogicalRelation => logicalRel.output.map(_.name).toSet
  case unaryNode: UnaryNode => findUsedInputColumns(unaryNode.child)
  case binaryNode: BinaryNode => findUsedInputColumns(binaryNode.left) ++ findUsedInputColumns(binaryNode.right)
  case otherNode => otherNode.children.flatMap(findUsedInputColumns).toSet
}

val usedColumns = findUsedInputColumns(outputDf.queryExecution.optimizedPlan)
println("实际使用的输入列:" + usedColumns.mkString(", "))

方法3:使用Spark 3.x+内置的SchemaPruning工具类

Spark 3.x提供了官方工具类用于裁剪Schema,可以直接提取查询所需的输入列:

import org.apache.spark.sql.catalyst.util.SchemaPruning

val inputPlan = inputDf.queryExecution.analyzed
val requiredColumns = SchemaPruning.pruneSchema(inputPlan, outputDf.queryExecution.analyzed)
  .map(_.name)

println("实际使用的输入列:" + requiredColumns.mkString(", "))

与解析explain方法的对比

上述方法均无需触发任何Action(如count()、show()),也避免了文本解析执行计划带来的出错风险,完全基于Spark内部API操作逻辑计划,性能和准确性更有保障。


内容的提问来源于stack exchange,提问作者Mehdi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 23:02:39