如何从Spark DataFrame中选取与case class字段对应的列
Spark Scala 通过Case Class自动选择DataFrame对应列的实现方法
核心思路是利用Spark原生的Encoder机制获取case class的字段列表,无需手动硬编码字段名,实现如下:
最简实现
import org.apache.spark.sql.Encoders import org.apache.spark.sql.functions.col val resultDf = df.select(Encoders.product[IdAndLastName].schema.fieldNames.map(col): _*)
实现原理
- case class默认继承
Product特质,Encoders.product[T]会根据case class的定义自动生成对应的Spark SQL结构信息 - 从生成的Schema中直接提取
fieldNames就能拿到和case class字段完全一致的字段名数组,严格按照你定义的case class字段顺序排列 - 用
: _*语法把数组展开作为select方法的可变参数,效果和手动写df.select("id","last")完全一致
通用优化(可选)
如果需要频繁做这类操作,可以给DataFrame扩展一个通用方法,后续调用更简洁:
import org.apache.spark.sql.DataFrame import org.apache.spark.sql.Encoder implicit class DFSelectExtension(df: DataFrame) { def selectCaseClassFields[T <: Product : Encoder]: DataFrame = { df.select(implicitly[Encoder[T]].schema.fieldNames.map(col): _*) } } // 调用只需一行,完全不需要写字段名 val resultDf = df.selectCaseClassFields[IdAndLastName]
注意事项
- 需保证case class的字段名和DataFrame中的列名大小写完全匹配
- 该方式会严格按照case class的字段顺序返回列,和手动指定顺序的效果一致
内容的提问来源于stack exchange,提问作者seriously divergent
相关产品推荐
相关产品推荐

