【问题标题】:Is there a better way to go about this process of trimming my spark DataFrame appropriately?有没有更好的方法来适当地修剪我的 spark DataFrame 的过程?
【发布时间】:2019-07-01 22:05:05
【问题描述】:

在下面的示例中,我希望能够只获取计数最高的 x 个 ID。 x 是我想要的这些数量,它由一个名为 howMany 的变量确定。

对于以下示例,给定此 Dataframe:

+------+--+-----+
|query |Id|count|
+------+--+-----+
|query1|11|2    |
|query1|12|1    |
|query2|13|2    |
|query2|14|1    |
|query3|13|2    |
|query4|12|1    |
|query4|11|1    |
|query5|12|1    |
|query5|11|2    |
|query5|14|1    |
|query5|13|3    |
|query6|15|2    |
|query6|16|1    |
|query7|17|1    |
|query8|18|2    |
|query8|13|3    |
|query8|12|1    |
+------+--+-----+

如果变量号为2,我想得到以下数据框。

+------+-------+-----+
|query |Ids    |count|
+------+-------+-----+
|query1|[11,12]|2    |
|query2|[13,14]|2    |
|query3|[13]   |2    |
|query4|[12,11]|1    |
|query5|[11,13]|2    |
|query6|[15,16]|2    |
|query7|[17]   |1    |
|query8|[18,13]|2    |
+------+-------+-----+

然后我想删除计数列,但这很简单。

我有办法做到这一点,但我认为它完全违背了 scala 的目的,并且完全浪费了大量的运行时间。作为新手,我不确定解决此问题的最佳方法

我目前的方法是首先获取查询列的不同列表并创建一个迭代器。其次,我使用迭代器遍历列表,并使用 df.select($"eachColumnName"...).where("query".equalTo(iter.next())) 将数据框修剪为列表中的当前查询.然后我 .limit(howMany) 和 groupBy($"query").agg(collect_list($"Id").as("Ids"))。最后,我有一个空数据框,并将这些中的每一个一个一个添加到空数据框并返回这个新创建的数据框。

df.select($"query").distinct().rdd.map(r => r(0).asInstanceOf[String]).collect().toList
val iter = queries.toIterator
while (iter.hasNext) {
    middleDF = df.select($"query", $"Id", $"count").where($"query".equalTo(iter.next()))
    queryDF = middleDF.sort(col("count").desc).limit(howMany).select(col("query"), col("Ids")).groupBy(col("query")).agg(collect_list("Id").as("Ids"))
    emptyDF.union(queryDF) // Assuming emptyDF is made
}
emptyDF

【问题讨论】:

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


    【解决方案1】:

    我会使用 Window-Functions 来获得排名,然后使用 groupBy 来聚合:

    import org.apache.spark.sql.expressions.Window
    import org.apache.spark.sql.functions._
    
    val howMany = 2
    
    val newDF = df
    .withColumn("rank",row_number().over(Window.partitionBy($"query").orderBy($"count".desc)))
    .where($"rank"<=howMany)
    .groupBy($"query")
    .agg(
     collect_list($"Id").as("Ids"),
     max($"count").as("count") 
    )
    

    【讨论】:

    • 啊,谢谢!那个等级库看起来很强大。我希望我自己找到了。非常感谢!
    猜你喜欢
    • 2015-07-27
    • 1970-01-01
    • 1970-01-01
    • 2015-07-31
    • 1970-01-01
    • 2018-12-12
    • 1970-01-01
    • 1970-01-01
    • 2010-09-14
    相关资源
    最近更新 更多