【问题标题】:How to merge map column in spark sql?如何在 spark sql 中合并地图列?
【发布时间】:2020-11-21 18:20:51
【问题描述】:

我在 Dataframe 中有两个 Map 类型的列。有没有一种方法可以创建一个新的 Map 列,使用 .withColumn 在 spark Sql 中合并这两列?

val sampleDF = Seq(
 ("Jeff", Map("key1" -> "val1"), Map("key2" -> "val2"))
).toDF("name", "mapCol1", "mapCol2")

sampleDF.show()

+----+-----------------+-----------------+
|name|          mapCol1|          mapCol2|
+----+-----------------+-----------------+
|Jeff|Map(key1 -> val1)|Map(key2 -> val2)|
+----+-----------------+-----------------+

【问题讨论】:

    标签: apache-spark apache-spark-sql


    【解决方案1】:

    您可以使用withColumn 编写一个udf 函数将两列合并为一个,如下所示

    import org.apache.spark.sql.functions._
    def mergeUdf = udf((map1: Map[String, String], map2: Map[String, String])=> map1 ++ map2)
    
    sampleDF.withColumn("merged", mergeUdf(col("mapCol1"), col("mapCol2"))).show(false)
    

    这应该给你

    +----+-----------------+-----------------+-------------------------------+
    |name|mapCol1          |mapCol2          |merged                         |
    +----+-----------------+-----------------+-------------------------------+
    |Jeff|Map(key1 -> val1)|Map(key2 -> val2)|Map(key1 -> val1, key2 -> val2)|
    +----+-----------------+-----------------+-------------------------------+
    

    希望回答对你有帮助

    【讨论】:

    • 谢谢!!这可行,但有没有不使用 udf 的方法?
    • 您可以使用数组或结构内置函数,但我认为您不想要结果
    • @Nats:现在可以使用map_concat 检查我对这个问题的回答。
    【解决方案2】:

    仅当由于性能原因您的用例没有内置函数时才使用 UDF。

    Spark 2.4 及以上版本

    import org.apache.spark.sql.functions.{map_concat, col}
    
    sampleDF.withColumn("map_concat", map_concat(col("mapCol1"), col("mapCol2"))).show(false)
    

    输出

    +----+-----------------+-----------------+-------------------------------+
    |name|mapCol1          |mapCol2          |map_concat                     |
    +----+-----------------+-----------------+-------------------------------+
    |Jeff|Map(key1 -> val1)|Map(key2 -> val2)|Map(key1 -> val1, key2 -> val2)|
    +----+-----------------+-----------------+-------------------------------+
    

    Spark 版本 2.4 以下

    按照 @RameshMaharjan answer in this question 创建一个 UDF,但我添加了一个空检查以避免在运行时出现 NPE,如果不添加最终会导致作业失败。

    import org.apache.spark.sql.functions.{udf, col}
    
    val map_concat = udf((map1: Map[String, String],
                          map2: Map[String, String]) =>
      if (map1 == null) {
        map2
      } else if (map2 == null) {
        map1
      } else {
        map1 ++ map2
      })
    
    sampleDF.withColumn("map_concat", map_concat(col("mapCol1"), col("mapCol2")))
     .show(false)
    

    【讨论】:

      【解决方案3】:

      你可以使用struct来实现。

      val sampleDF = Seq(
       ("Jeff", Map("key1" -> "val1"), Map("key2" -> "val2"))
      ).toDF("name", "mapCol1", "mapCol2")
      
      sampleDF.show()
      
      +----+-----------------+-----------------+
      |name|          mapCol1|          mapCol2|
      +----+-----------------+-----------------+
      |Jeff|Map(key1 -> val1)|Map(key2 -> val2)|
      +----+-----------------+-----------------+
      
      sampleDF.withColumn("NewColumn",struct(sampleDF("mapCol1"), sampleDF("mapCol2"))).take(2)
          res17: Array[org.apache.spark.sql.Row] = Array([Jeff,Map(key1 -> val1),Map(key2 -> val2),[Map(key1 -> val1),Map(key2 -> val2)]])
      
      +----+-----------------+-----------------+--------------------+
      |name|          mapCol1|          mapCol2|           NewColumn|
      +----+-----------------+-----------------+--------------------+
      |Jeff|Map(key1 -> val1)|Map(key2 -> val2)|[Map(key1 -> val1...|
      +----+-----------------+-----------------+--------------------+
      

      参考:How to merge two columns of a `Dataframe` in Spark into one 2-Tuple?

      【讨论】:

      • 这不会合并地图,它会创建一个具有 2 个地图字段的结构
      猜你喜欢
      • 2020-07-17
      • 1970-01-01
      • 2023-04-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-05-02
      • 1970-01-01
      相关资源
      最近更新 更多