【问题标题】:create empty array-column of given schema in Spark在 Spark 中创建给定模式的空数组列
【发布时间】:2018-12-06 04:53:19
【问题描述】:

由于 parquet 无法解析空数组,我在写表之前将空数组替换为 null。现在,当我阅读表格时,我想做相反的事情:

我有一个具有以下架构的 DataFrame:

|-- id: long (nullable = false)
 |-- arr: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- x: double (nullable = true)
 |    |    |-- y: double (nullable = true)

以及以下内容:

+---+-----------+
| id|        arr|
+---+-----------+
|  1|[[1.0,2.0]]|
|  2|       null|
+---+-----------+

我想将空数组 (id=2) 替换为空数组,即

+---+-----------+
| id|        arr|
+---+-----------+
|  1|[[1.0,2.0]]|
|  2|         []|
+---+-----------+

我试过了:

val arrSchema = df.schema(1).dataType

df
.withColumn("arr",when($"arr".isNull,array().cast(arrSchema)).otherwise($"arr"))
.show()

给出:

java.lang.ClassCastException: org.apache.spark.sql.types.NullType$ 无法转换为 org.apache.spark.sql.types.StructType

编辑:我不想“硬编码”我的数组列的任何架构(至少不是结构的架构),因为这可能因情况而异。我只能在运行时使用来自df 的架构信息

顺便说一句,我使用的是 Spark 2.1,因此我无法使用 typedLit

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:
    • 具有已知外部类型的 Spark 2.2+

      一般你可以使用typedLit来提供空数组。

      import org.apache.spark.sql.functions.typedLit
      
      typedLit(Seq.empty[(Double, Double)])
      

      要为嵌套对象使用特定名称,您可以使用案例类:

      case class Item(x: Double, y: Double)
      
      typedLit(Seq.empty[Item])
      

      rename by cast:

      typedLit(Seq.empty[(Double, Double)])
        .cast("array<struct<x: Double, y: Double>>")
      
    • Spark 2.1+ 仅带有架构

      只有架构你可以尝试:

      val schema = StructType(Seq(
        StructField("arr", StructType(Seq(
          StructField("x", DoubleType),
          StructField("y", DoubleType)
        )))
      ))
      
      def arrayOfSchema(schema: StructType) =
        from_json(lit("""{"arr": []}"""), schema)("arr")
      
      arrayOfSchema(schema).alias("arr")
      

      其中schema 可以从现有的DataFrame 中提取并用额外的StructType 包装:

      StructType(Seq(
        StructField("arr", df.schema("arr").dataType)
      ))
      

    【讨论】:

    • 是否可以只使用来自 df 的模式信息(这会使其更通用)
    • 其实我用UDF找到了一个更简单的解决方案
    • @RaphaelRoth Neat.
    【解决方案2】:

    一种方法是使用 UDF:

    val arrSchema = df.schema(1).dataType // ArrayType(StructType(StructField(x,DoubleType,true), StructField(y,DoubleType,true)),true)
    
    val emptyArr = udf(() => Seq.empty[Any],arrSchema)
    
    df
    .withColumn("arr",when($"arr".isNull,emptyArr()).otherwise($"arr"))
    .show()
    
    +---+-----------+
    | id|        arr|
    +---+-----------+
    |  1|[[1.0,2.0]]|
    |  2|         []|
    +---+-----------+
    

    【讨论】:

      【解决方案3】:

      另一种方法是使用coalesce:

      val df = Seq(
        (Some(1), Some(Array((1.0, 2.0)))),
        (Some(2), None)
      ).toDF("id", "arr")
      
      df.withColumn("arr", coalesce($"arr", typedLit(Array.empty[(Double, Double)]))).
        show
      // +---+-----------+
      // | id|        arr|
      // +---+-----------+
      // |  1|[[1.0,2.0]]|
      // |  2|         []|
      // +---+-----------+
      

      【讨论】:

      • 这也仅适用于 Spark 2.2+,并且确实需要硬编码的类型信息
      【解决方案4】:

      带有案例类的 UDF 也可能很有趣:

      case class Item(x: Double, y: Double)
      val udf_emptyArr = udf(() => Seq[Item]())
      df
      .withColumn("arr",coalesce($"arr",udf_emptyArr()))
      .show()
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2016-04-28
        • 1970-01-01
        • 1970-01-01
        • 2018-02-07
        • 2020-01-01
        • 2019-11-26
        相关资源
        最近更新 更多