【问题标题】:Extracting `Seq[(String,String,String)]` from spark DataFrame从 spark DataFrame 中提取 `Seq[(String,String,String)]`
【发布时间】:2016-05-31 18:25:48
【问题描述】:

我有一个带有 Seq[(String, String, String)] 行的 spark DF。我正在尝试对此进行某种flatMap,但我所做的任何尝试最终都会抛出

java.lang.ClassCastException:org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema 无法转换为 scala.Tuple3

我可以从 DF 中取出单行或多行

df.map{ r => r.getSeq[Feature](1)}.first

返回

Seq[(String, String, String)] = WrappedArray([ancient,jj,o], [olympia_greece,nn,location] .....

并且 RDD 的数据类型似乎是正确的。

org.apache.spark.rdd.RDD[Seq[(String, String, String)]]

df 的架构是

root
 |-- article_id: long (nullable = true)
 |-- content_processed: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- lemma: string (nullable = true)
 |    |    |-- pos_tag: string (nullable = true)
 |    |    |-- ne_tag: string (nullable = true)

我知道这个问题与 spark sql 将 RDD 行视为org.apache.spark.sql.Row 有关,即使他们愚蠢地说这是Seq[(String, String, String)]。有一个相关的问题(下面的链接),但该问题的答案对我不起作用。我对 spark 还不够熟悉,无法弄清楚如何将其转变为可行的解决方案。

这些行是Row[Seq[(String, String, String)]] 还是Row[(String, String, String)] 还是Seq[Row[(String, String, String)]] 或者更疯狂的东西。

我正在尝试做类似的事情

df.map{ r => r.getSeq[Feature](1)}.map(_(1)._1)

看似有效,但实际上无效

df.map{ r => r.getSeq[Feature](1)}.map(_(1)._1).first

抛出上述错误。那么我应该如何(例如)获取每行第二个元组的第一个元素?

另外为什么已经设计了 spark 来做到这一点,声称某物是一种类型而实际上它不是并且不能转换为声称的类型似乎是愚蠢的。 p>


相关问题:GenericRowWithSchema exception in casting ArrayBuffer to HashSet in DataFrame to RDD from Hive table

相关错误报告:http://search-hadoop.com/m/q3RTt2bvwy19Dxuq1&subj=ClassCastException+when+extracting+and+collecting+DF+array+column+type

【问题讨论】:

  • 请问谁否决了这个问题以及为什么?
  • 相关的错误报告是我的解决方案

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


【解决方案1】:

好吧,它并没有声称它是一个元组。它声称它是一个struct,它映射到Row

import org.apache.spark.sql.Row

case class Feature(lemma: String, pos_tag: String, ne_tag: String)
case class Record(id: Long, content_processed: Seq[Feature])

val df = Seq(
  Record(1L, Seq(
    Feature("ancient", "jj", "o"),
    Feature("olympia_greece", "nn", "location")
  ))
).toDF

val content = df.select($"content_processed").rdd.map(_.getSeq[Row](0))

您可以在Spark SQL programming guide 中找到确切的映射规则。

由于Row 的结构并不完全漂亮,您可能希望将其映射到有用的地方:

content.map(_.map {
  case Row(lemma: String, pos_tag: String, ne_tag: String) => 
    (lemma, pos_tag, ne_tag)
})

或:

content.map(_.map ( row => (
  row.getAs[String]("lemma"),
  row.getAs[String]("pos_tag"),
  row.getAs[String]("ne_tag")
)))

最后一个更简洁的方法是Datasets

df.as[Record].rdd.map(_.content_processed)

df.select($"content_processed").as[Seq[(String, String, String)]]

虽然此时这似乎有点错误。

第一种方法 (Row.getAs) 和第二种方法 (Dataset.as) 有重要区别。前者将对象提取为Any 并应用asInstanceOf。后一种是使用编码器在内部类型和所需表示之间进行转换。

【讨论】:

  • 是的,我没有意识到我可以将内容作为基本上任何内容来询问,并且只会在运行时进行检查(我确信在文档中的某处提到了这一点)。 df.select($"content_processed").map(_.getSeq[(String, String)](0)) 也返回 something,但实际上不会运行。我假设这不是从 Cassandra 中提取数据的人工制品,而是 DataFrames 固有的东西?
  • 嗯,是的。 Row 只是 Any 的集合。在getSeq 和其他方法中发生的所有事情都等同于(row(i): Any).asInstanceOf[T]
【解决方案2】:
object ListSerdeTest extends App {

  implicit val spark: SparkSession = SparkSession
    .builder
    .master("local[2]")
    .getOrCreate()


  import spark.implicits._
  val myDS = spark.createDataset(
    Seq(
      MyCaseClass(mylist = Array(("asd", "aa"), ("dd", "ee")))
    )
  )

  myDS.toDF().printSchema()

  myDS.toDF().foreach(
    row => {
      row.getSeq[Row](row.fieldIndex("mylist"))
        .foreach {
          case Row(a, b) => println(a, b)
        }
    }
  )
}

case class MyCaseClass (
                 mylist: Seq[(String, String)]
               )

以上代码是处理嵌套结构的另一种方法。 Spark 默认编码器将对 TupleX 进行编码,使它们成为嵌套结构,这就是您看到这种奇怪行为的原因。就像其他人在评论中所说的那样,你不能只做getAs[T](),因为它只是一个演员(x.asInstanceOf[T]),因此会给你运行时异常。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2022-08-23
    • 1970-01-01
    • 2022-07-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-28
    相关资源
    最近更新 更多