【问题标题】:Spark Scala groupBy multiple columns with valuesSpark Scala groupBy 具有值的多列
【发布时间】:2020-03-09 21:15:03
【问题描述】:

我在 spark 中有一个以下数据框 (df)

| group_1 | group_2 | year | value |
| "School1" | "Student" | 2018 | name_aaa |
| "School1" | "Student" | 2018 | name_bbb |
| "School1" | "Student" | 2019 | name_aaa |
| "School2" | "Student" | 2019 | name_aaa |

我想要的是

| group_1 | group_2 | values_map |
| "School1" | "Student" | [2018 -> [name_aaa, name_bbb], [2019 -> [name_aaa] |
| "School2" | "Student" | [2019 -> [name_aaa] |

我用groupBy 和collect_list() 和map() 尝试过,但没有成功。它创建了一个只有来自name_aaa 或name_bbb 的最后一个值的映射。如何使用 Apache Spark 实现这一目标?

【问题讨论】:

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


    【解决方案1】:

    另一个答案的结果是数组类型而不是映射。这是为您的结果实现map 类型列的方法。

    df.groupBy("group_1", "group_2", "year").agg(collect_list("value").as("value_list"))
      .groupBy("group_1", "group_2").agg(collect_list(struct(col("year"), col("value_list"))).as("map_list"))
      .withColumn("values_map", map_from_entries(col("map_list")))
      .drop("map_list")
      .show(false)
    

    我没有使用过udf。然后,结果直接显示您的预期。

    +-------+-------+--------------------------------------------------+
    |group_1|group_2|values_map                                        |
    +-------+-------+--------------------------------------------------+
    |School2|Student|[2019 -> [name_aaa]]                              |
    |School1|Student|[2018 -> [name_aaa, name_bbb], 2019 -> [name_aaa]]|
    +-------+-------+--------------------------------------------------+
    

    【讨论】:

      【解决方案2】:

      解决方案可能是:

      scala> df1.show
      +-------+-------+----+--------+
      |group_1|group_2|year|   value|
      +-------+-------+----+--------+
      |school1|student|2018|name_aaa|
      |school1|student|2018|name_bbb|
      |school1|student|2019|name_aaa|
      |school2|student|2019|name_aaa|
      +-------+-------+----+--------+
      
      
      scala> val df2 = df1.groupBy("group_1","group_2","year").agg(collect_list('value).as("value"))
      df2: org.apache.spark.sql.DataFrame = [group_1: string, group_2: string ... 2 more fields]
      
      scala> df2.show
      +-------+-------+----+--------------------+
      |group_1|group_2|year|               value|
      +-------+-------+----+--------------------+
      |school1|student|2018|[name_aaa, name_bbb]|
      |school1|student|2019|          [name_aaa]|
      |school2|student|2019|          [name_aaa]|
      +-------+-------+----+--------------------+
      
      
      scala> val myUdf = udf((year: String, values: Seq[String]) => Map(year -> values))
      myUdf: org.apache.spark.sql.expressions.UserDefinedFunction = UserDefinedFunction(<function2>,MapType(StringType,ArrayType(StringType,true),true),Some(List(StringType, ArrayType(StringType,true))))
      
      scala> val df3 = df2.withColumn("values",myUdf($"year",$"value")).drop("year","value")
      df3: org.apache.spark.sql.DataFrame = [group_1: string, group_2: string ... 1 more field]
      scala> val df4 = df3.groupBy("group_1","group_2").agg(collect_list("values").as("value_map"))
      df4: org.apache.spark.sql.DataFrame = [group_1: string, group_2: string ... 1 more field]
      
      scala> df4.printSchema
      root
       |-- group_1: string (nullable = true)
       |-- group_2: string (nullable = true)
       |-- value_map: array (nullable = true)
       |    |-- element: map (containsNull = true)
       |    |    |-- key: string
       |    |    |-- value: array (valueContainsNull = true)
       |    |    |    |-- element: string (containsNull = true)
      
      
      scala> df4.show(false)
      +-------+-------+------------------------------------------------------+
      |group_1|group_2|value_map                                             |
      +-------+-------+------------------------------------------------------+
      |school1|student|[[2018 -> [name_aaa, name_bbb]], [2019 -> [name_aaa]]]|
      |school2|student|[[2019 -> [name_aaa]]]                                |
      +-------+-------+------------------------------------------------------+
      

      如果有帮助请告诉我!!

      【讨论】:

      • 你能告诉我投反对票的原因吗?反馈将不胜感激。
      • 谢谢,它有效。如果我有Struct&lt;...&gt; 而不是String,我应该用什么更改Seq[String]?编辑:我已经赞成/接受您的解决方案。可能是其他人否决了您的答案。
      • 你能给我架构的详细信息吗?
      • 是人(id: Int, firstName: String, lastName: String)
      • 谢谢哥们!!我认为这个链接:stackoverflow.com/question/42931796/… 可能会帮助您找到如何将struct 类型传递给 udf。如果有任何问题,请告诉我!
      猜你喜欢
      • 2018-09-09
      • 1970-01-01
      • 2018-08-20
      • 2021-08-30
      • 2021-12-02
      • 1970-01-01
      • 1970-01-01
      • 2019-05-25
      • 1970-01-01
      相关资源
      最近更新 更多