【问题标题】:How to convert Row of a Scala DataFrame into case class most efficiently?如何最有效地将 Scala DataFrame 的行转换为案例类?
【发布时间】:2015-03-25 20:03:58
【问题描述】:

在 Spark 中获得一些 Row 类(Dataframe 或 Catalyst)后,我想在我的代码中将其转换为案例类。这可以通过匹配来完成

someRow match {case Row(a:Long,b:String,c:Double) => myCaseClass(a,b,c)}

但是当行有大量列时,它会变得很难看,比如十几个双精度数、一些布尔值甚至偶尔的空值。

我希望能够 - 抱歉 - 将 Row 转换为 myCaseClass。有没有可能,或者我已经得到了最经济的语法?

【问题讨论】:

  • 可能无形(github.com/milessabin/shapeless/wiki/…)可以帮助减少样板,但可能不太喜欢nulls。也许宏(如果你有很多案例类)?
  • 从未尝试过宏。这里的一个问题是我是语言标准的信徒。我可以想象我总是可以做我自己的方法,或者使用其他人......但我更愿意尝试理解它是如何在没有任何外部因素的情况下完成的。
  • 想知道...也许我可以从 Row 继承“myCaseClass”?
  • 这太令人失望了。我有一个大而复杂的案例类,现在需要在我想加载和使用它时手动将每一列映射回它。这让我很难过:-(

标签: scala apache-spark apache-spark-sql


【解决方案1】:

DataFrame 只是 Dataset[Row] 的类型别名。与强类型化 Scala/Java 数据集附带的“类型化转换”相比,这些操作也称为“非类型化转换”。

在 spark 中从 Dataset[Row] 到 Dataset[Person] 的转换非常简单

val DFtoProcess = SQLContext.sql("SELECT * FROM peoples WHERE name='test'")

此时,Spark 会将您的数据转换为 DataFrame = Dataset[Row],这是一个通用 Row 对象的集合,因为它不知道确切的类型。

// Create an Encoders for Java class (In my eg. Person is a JAVA class)
// For scala case class you can pass Person without .class reference
val personEncoder = Encoders.bean(Person.class) 

val DStoProcess = DFtoProcess.as[Person](personEncoder)

现在,Spark 转换 Dataset[Row] -> Dataset[Person] 类型特定的 Scala / Java JVM 对象,由 Person 类指定。

详情请参考以下databricks提供的链接

https://databricks.com/blog/2016/07/14/a-tale-of-three-apache-spark-apis-rdds-dataframes-and-datasets.html

【讨论】:

  • 您只有一个答案 - 但无论如何它都是一个好答案!在偶然发现这个答案之前,我找不到关于如何创建自定义 Spark 编码器的 任何 信息。顺便说一句scala 方式是Encoders.bean[Person]
  • 应该注意as 是不安全的,因为它不检查强制转换是否有效。我不明白他们是如何在他们的 API 中添加如此丑陋的功能的(它应该被称为 unsafeAs - 同时进行过滤或返回 DataSet[Option[T]] 的安全 as 在哪里?)
  • @javadba - 我猜这不适用于案例类?我在这里收到错误Cannot infer type for class Person because it is not bean-compliant
  • 当我想将Dataset[Row] 转换为Dataset[Person] 时,这显然是最干净的解决方案,但是如果我只想转换单个Row 对象怎么办?我只能从Row 中的每个字段手动构建Person。有没有更好的办法?
  • @sasgorilla 案例类的正确答案如下 - 你需要 import spark.implicits._ 然后你可以使用 df.as[Person]
【解决方案2】:

据我所知,您不能将 Row 强制转换为案例类,但我有时会选择直接访问行字段,例如

map(row => myCaseClass(row.getLong(0), row.getString(1), row.getDouble(2))

我发现这更容易,特别是如果案例类构造函数只需要行中的一些字段。

【讨论】:

  • 你避免了匹配java nulls的问题:-)
  • 对于较小的列集,我也喜欢这种表示形式,但如果列集较大会增加歧义,那么我认为@Gianmarios 的建议可能更具可扩展性。我需要自己验证几件事。会就此与您联系。
  • 如果案例类中的某些字段是通用的,那会起作用吗?
  • 谢谢!这是一个很大的帮助。
【解决方案3】:
scala> import spark.implicits._    
scala> val df = Seq((1, "james"), (2, "tony")).toDF("id", "name")
df: org.apache.spark.sql.DataFrame = [id: int, name: string]

scala> case class Student(id: Int, name: String)
defined class Student

scala> df.as[Student].collectAsList
res6: java.util.List[Student] = [Student(1,james), Student(2,tony)]

这里spark.implicits._ 中的spark 是您的SparkSession。如果您在 REPL 中,则会话已定义为 spark,否则您需要相应地调整名称以对应您的 SparkSession

【讨论】:

  • 对于 Spark 2.1.0,我必须导入 spark.implicits._ 才能完成这项工作 - Scala 的漂亮、优雅的解决方案
  • 有关spark.implicits._的更多详细信息,请参阅stackoverflow.com/questions/39968707/…
  • 我在 REPL 之外和测试中执行此操作,它显示 Error:(44, 30) not enough arguments for method as: (implicit evidence$2: org.apache.spark.sql.Encoder[Student])org.apache.spark.sql.Dataset[Student]. Unspecified value parameter evidence$2. val rows = df.as[Student].collectAsList()
【解决方案4】:

当然,您可以将 Row 对象匹配到案例类中。假设您的 SchemaType 有许多字段,并且您希望将其中一些字段匹配到您的案例类中。 如果你没有空字段,你可以简单地做:

case class MyClass(a: Long, b: String, c: Int, d: String, e: String)

dataframe.map {
  case Row(a: java.math.BigDecimal, 
    b: String, 
    c: Int, 
    _: String,
    _: java.sql.Date, 
    e: java.sql.Date,
    _: java.sql.Timestamp, 
    _: java.sql.Timestamp, 
    _: java.math.BigDecimal, 
    _: String) => MyClass(a = a.longValue(), b = b, c = c, d = d.toString, e = e.toString)
}

这种方法在空值的情况下会失败,并且还需要您明确定义每个字段的类型。 如果您必须处理空值,则需要丢弃所有包含空值的行

dataframe.na.drop()

即使空字段不是在您的案例类的模式匹配中使用的字段,也会删除记录。 或者,如果您想处理它,您可以将 Row 对象转换为 List,然后使用选项模式:

case class MyClass(a: Long, b: String, c: Option[Int], d: String, e: String)

dataframe.map(_.toSeq.toList match {
  case List(a: java.math.BigDecimal, 
    b: String, 
    c: Int, 
    _: String,
    _: java.sql.Date, 
    e: java.sql.Date,
    _: java.sql.Timestamp, 
    _: java.sql.Timestamp, 
    _: java.math.BigDecimal, 
    _: String) => MyClass(
      a = a.longValue(), b = b, c = Option(c), d = d.toString, e = e.toString)
}

查看这个 github 项目 Sparkz(),它将很快引入许多用于简化 Spark 和 DataFrame API 并使它们更加面向函数式编程的库。

【讨论】:

猜你喜欢
  • 2014-01-08
  • 1970-01-01
  • 2016-07-27
  • 2014-12-18
  • 1970-01-01
  • 2015-06-16
  • 1970-01-01
  • 2016-08-28
  • 2015-10-23
相关资源
最近更新 更多