Spark中如何检测输入DataFrame里被转换逻辑实际使用的列?
检测Spark DataFrame中被实际使用的输入列
方法1:直接解析分析后的逻辑计划(推荐)
不需要触发Action,也不用解析文本格式的执行计划,直接通过操作Spark内部的逻辑计划结构提取被引用的列,可靠性和效率更高:
- 获取输出DataFrame经过分析后的逻辑计划
- 递归遍历计划节点,收集所有被引用的
Attribute - 关联这些属性到输入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
相关产品推荐
相关产品推荐

