【问题标题】:Spark dataframe grouping, sorting, and selecting top rows for a set of columnsSpark 数据框分组、排序和选择一组列的顶行
【发布时间】:2019-01-13 00:04:40
【问题描述】:

我使用的是 Spark 1.5.0。我有一个包含以下列的 Spark 数据框:

| user_id | description | fName | weight |

我想做的是为每个用户选择前 10 行和后 10 行(基于列权重的值,数据类型为 Double)。如何使用 Spark SQL 或数据帧操作来做到这一点?

例如。为简单起见,我只为每个用户选择前 2 行(基于重量)。我想根据绝对重量的值对o/p进行排序。

u1  desc1 f1  -0.20
u1  desc1 f1  +0.20
u2  desc1 f1  0.80
u2  desc1 f1  -0.60
u1  desc1 f1  1.10
u1  desc1 f1  6.40
u2  desc1 f1  0.05
u1  desc1 f1  -3.20
u2  desc1 f1  0.50
u2  desc1 f1  -0.70
u2  desc1 f1  -0.80   

这是所需的o/p:

u1  desc1 f1  6.40
u1  desc1 f1  -3.20
u1  desc1 f1  1.10
u1  desc1 f1  -0.20
u2  desc1 f1  0.80
u2  desc1 f1  -0.80
u2  desc1 f1  -0.70
u2  desc1 f1  0.50

【问题讨论】:

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


    【解决方案1】:

    您可以通过row_number 使用窗口函数:

    import org.apache.spark.sql.functions.row_number
    import org.apache.spark.sql.expressions.Window
    
    val w = Window.partitionBy($"user_id")
    val rankAsc = row_number().over(w.orderBy($"weight")).alias("rank_asc")
    val rankDesc = row_number().over(w.orderBy($"weight".desc)).alias("rank_desc")
    
    df.select($"*", rankAsc, rankDesc).filter($"rank_asc" <= 2 || $"rank_desc" <= 2)
    

    在 Spark 1.5.0 中,您可以使用 rowNumber 代替 row_number

    【讨论】:

    • 我认为 row_number 不包含在 Spark 1.5 中。我收到此错误:scala> import org.apache.spark.sql.functions.row_number :29: error: value row_number is not a member of object org.apache.spark.sql.functions import org.apache.spark .sql.functions.row_number
    • 你是对的,但即使函数名称不同,解决方案也是相同的 (rowNumber)。
    • 行号(单独计算)如何与数据集 df 记录关联?
    • @SumitKumarGhosh 我不确定我是否跟随。你问rankAsc/rankDesc吗?那里没有计算。这些只是用于为filter 生成执行计划的 SQL 表达式。
    猜你喜欢
    • 2018-10-24
    • 2020-08-30
    • 2017-06-30
    • 1970-01-01
    • 1970-01-01
    • 2020-04-12
    • 1970-01-01
    • 2016-06-28
    相关资源
    最近更新 更多