【问题标题】:Convert Dataset<Row> to a Typed Dataset having optional Parameters将 Dataset<Row> 转换为具有可选参数的类型化数据集
【发布时间】:2020-05-20 12:19:20
【问题描述】:

我有一个数据集,我希望将其转换为类型数据集,其中类型是具有 Option 多个参数的案例类。例如,我使用 spark shell 创建了一个案例类、一个编码器和(原始)数据集:

case class Analogue(id: Long, t1: Option[Double] = None, t2: Option[Double] = None)
val df = Seq((1, 34.0), (2,3.4)).toDF("id", "t1")
implicit val analogueChannelEncoder: Encoder[Analogue] = Encoders.product[Analogue]

我想从df 创建一个Dataset&lt;Analogue&gt;,所以我尝试:

df.as(analogueChannelEncoder)

但这会导致错误:

org.apache.spark.sql.AnalysisException: cannot resolve '`t2`' given input columns: [id, t1];

查看dfanalogueChannelEncoder 的架构,区别很明显:

scala> df.schema
res3: org.apache.spark.sql.types.StructType = StructType(StructField(id,IntegerType,false), StructField(t1,DoubleType,false))

scala> analogueChannelEncoder.schema
res4: org.apache.spark.sql.types.StructType = StructType(StructField(id,LongType,false), StructField(t1,DoubleType,true), StructField(t2,DoubleType,true))

我已经看到 this 的答案,但这对我不起作用,因为我的 Dataset 已组装并且不是来自数据源的直接加载

如何将未键入的Dataset&lt;Row&gt; 转换为Dataset&lt;Analogue&gt;

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    你的案例类

      case class Analogue(id: Long, t1: Option[Double] = None, t2: Option[Double] = None)
    
    

    您的转换代码...

     val encoderSchema = Encoders.product[Analogue].schema
      val df1: Dataset[Row] = spark.createDataset(Seq((1, 34.0), (2, 3.4))).map(x => Analogue(x._1, Option(x._2), None))
        .toDF("id", "t1", "t2")
      df1.show
    
      df1.printSchema()
      encoderSchema.printTreeString()
    

    结果:

    +---+----+----+
    | id|  t1|  t2|
    +---+----+----+
    |  1|34.0|null|
    |  2| 3.4|null|
    +---+----+----+
    
    root
     |-- id: long (nullable = false)
     |-- t1: double (nullable = true)
     |-- t2: double (nullable = true)
    
    root
     |-- id: long (nullable = false)
     |-- t1: double (nullable = true)
     |-- t2: double (nullable = true)
    

    更新(在案例类中增加了更多作为选项的列):

    我假设您的案例类有很多字段(例如 5 个字段) 如果选项值为无,则... 它的工作原理如下...

    下面是例子

    
      case class Analogue(id: Long, t1: Option[Double] = None, t2: Option[Double] = None, t3: Option[Double] = None, t4: Option[Double] = None, t5: Option[Double] = None)
    
    
    
      val encoderSchema = Encoders.product[Analogue].schema
      println(encoderSchema.toSeq)
      val df1 = spark.createDataset(Seq((1, 34.0), (2, 3.4)))
        .map(x => Analogue(x._1, Option(x._2)))
         .as[Analogue].toDF()
      df1.show
      df1.printSchema()
      encoderSchema.printTreeString()
    
    

    如果您设置存在的字段,其余字段将被视为无。

    StructType(StructField(id,LongType,false), StructField(t1,DoubleType,true), StructField(t2,DoubleType,true), StructField(t3,DoubleType,true), StructField(t4,DoubleType,true), StructField(t5,DoubleType,true))
    +---+----+----+----+----+----+
    | id|  t1|  t2|  t3|  t4|  t5|
    +---+----+----+----+----+----+
    |  1|34.0|null|null|null|null|
    |  2| 3.4|null|null|null|null|
    +---+----+----+----+----+----+
    
    root
     |-- id: long (nullable = false)
     |-- t1: double (nullable = true)
     |-- t2: double (nullable = true)
     |-- t3: double (nullable = true)
     |-- t4: double (nullable = true)
     |-- t5: double (nullable = true)
    
    root
     |-- id: long (nullable = false)
     |-- t1: double (nullable = true)
     |-- t2: double (nullable = true)
     |-- t3: double (nullable = true)
     |-- t4: double (nullable = true)
     |-- t5: double (nullable = true)
    
    

    如果它不能以这种方式工作,请考虑我的评论广播想法并进一步工作。

    【讨论】:

    • 然而,这是一个很好的答案,这是我的错,在我的示例中,我没有明确表示可能指定了 t2 而未指定 t1。不幸的是,您的解决方案针对可用数据的一种变体进行了硬编码(在第二行)。在我的实际问题中,我有 50 个或更多可选变量,因此需要避免硬编码数据组合。
    • 我认为您可以广播encoderSchema.toSeq(这是目标模式),然后根据位置值,您可以将 Option 替换为 None 或任何相关的内容。为此,您需要使用val mytargetschema = spark.sparkContext.broadcast(encoderSchema.toSeq),您将获得像 mytargetschema.value 这样的地图(如上所示),然后您可以在没有硬编码的情况下使用它......需要做更多的努力。:-)..
    • @RamGhadiyaram.. 我在 spark 2.4 中收到此错误 java.lang.ClassCastException: $line9.$read$$iw$$iw$Analogue cannot be cast to $line9.$read$$iw$$iw$Analogue
    • 在我在这里添加的测试之后,上面的代码非常好。请检查你可能犯了一些愚蠢的错误
    • @Kerry :有用吗?
    【解决方案2】:

    我已通过检查“传入”Dataset&lt;Row&gt; 的列并将它们与Dataset&lt;Analogue&gt; 中的列进行比较来解决此问题。我用来将新列附加到 Dataset&lt;Row&gt; 然后将其转换为 Dataset&lt;Analogue&gt; 之前的结果差异。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-04-20
      • 1970-01-01
      • 1970-01-01
      • 2020-05-23
      • 2020-10-21
      • 2019-12-04
      • 2022-11-07
      相关资源
      最近更新 更多