【问题标题】:Spark dataframe aggregating the values聚合值的 Spark 数据框
【发布时间】:2020-02-26 15:06:00
【问题描述】:

数据框输入

+-----------------+-------+
|Id               | value |
+-----------------+-------+
|             1622| 139685|
|             1622| 182118|
|             1622| 127955|
|             3837|3224815|
|             1622| 727761|
|             1622| 155875|
|             3837|1504923|
|             1622| 139684|
|             1453| 536111|
+-----------------+-------+

输出:

    +-----------------+--------------------------------------------+
    |Id               | value                                      |
    +-----------------+--------------------------------------------+
    |             1622|[139685,182118,127955,727761,155875,139684] |
    |             1453| 536111                                     |
    |             3837|[3224815,1504923]                           |
    +-----------------+--------------------------------------------+


当特定的id 具有多个值时,它应该以array 格式收集,否则 它应该将其视为单个值withoutbracket []

我尝试使用以下链接解决方案,但无法处理数据框中的 if-else 条件。

链接:Spark DataFrame aggregate column values by key into List

【问题讨论】:

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


    【解决方案1】:

    使用窗口函数

    scala> import org.apache.spark.sql.expressions.Window
    scala> var df = Seq((1622, 139685),(1622, 182118),(1622, 127955),(3837,3224815),(1622, 727761),(1622, 155875),(3837,1504923),(1622, 139684),(1453, 536111)).toDF("id","value")
    
    scala> df.show()
    +----+-------+
    |  id|  value|
    +----+-------+
    |1622| 139685|
    |1622| 182118|
    |1622| 127955|
    |3837|3224815|
    |1622| 727761|
    |1622| 155875|
    |3837|1504923|
    |1622| 139684|
    |1453| 536111|
    +----+-------+
    
    scala> var df1= df.withColumn("r",count($"id").over(Window.partitionBy("id").orderBy("id")).cast("int"))
    
    scala> df1.show()
    +----+-------+---+
    |  id|  value|  r|
    +----+-------+---+
    |1453| 536111|  1|
    |1622| 139685|  6|
    |1622| 182118|  6|
    |1622| 127955|  6|
    |1622| 727761|  6|
    |1622| 155875|  6|
    |1622| 139684|  6|
    |3837|3224815|  2|
    |3837|1504923|  2|
    +----+-------+---+
    scala> var df2 =df1.selectExpr("*").filter('r ===1).drop("r").union(df1.filter('r =!= 1).groupBy("id").agg(collect_list($"value").cast("string").as("value")))
    
    
    scala> df2.show(false)
    +----+------------------------------------------------+
    |id  |value                                           |
    +----+------------------------------------------------+
    |1453|536111                                          |
    |1622|[139685, 182118, 127955, 727761, 155875, 139684]|
    |3837|[3224815, 1504923]                              |
    +----+------------------------------------------------+
    scala> df2.printSchema
    root
     |-- id: integer (nullable = false)
     |-- value: string (nullable = true)
    

    如果您有任何与此相关的问题,请告诉我。

    【讨论】:

    • 谢谢,我想你没有理解这个问题。如果id 有一个值,那么我不应该以列表格式显示。例如:1453 | [536111]
    • 我希望1453 | [536111] 应该是1453 | 536111 这种格式
    • @Manju,单列中不能有多种格式。在这里你将有一个ArrayType
    • 我们不能将这个值[123,1324]1234 存储在字符串格式中吗? @BlueSheepToken 而不是特定的数据类型格式更好地保留在字符串格式中
    • 你能详细说明一下df1.selectExpr("*").filter('r ===1).drop("r").union(df1.filter('r =!= 1).groupBy("id").agg(collect_list($"value").cast("string").as("value"))) 这个说法吗?
    猜你喜欢
    • 1970-01-01
    • 2015-05-19
    • 2020-11-20
    • 2019-02-27
    • 1970-01-01
    • 2017-05-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多