【问题标题】:Use Iterator to get top k keywords使用 Iterator 获取前 k 个关键字
【发布时间】:2020-02-08 09:15:19
【问题描述】:

我正在编写一个 Spark 算法来获取每个国家/地区的前 k 个关键字,现在我已经有一个包含所有记录的 Dataframe 并计划这样做

df.repartition($"country_id").mapPartition()

检索前 k 个关键字,但对如何编写迭代器来获取它感到困惑。

如果我能够编写方法或调用本机方法,我可以在每个分区中排序并获得前 k 个,如果输入是迭代器,这似乎不是正确的方法。

有人知道吗?

【问题讨论】:

  • 你能举个df的例子吗?
  • 不相信您需要 mapPartition
  • df 就像 country_id |关键字名称 | search_freq,如果我们已经repartitionBy country_id那么我们只需要在每个分区中对search_freq进行排序就可以得到前k个keyword_name,这样有意义吗?

标签: scala apache-spark pyspark iterator


【解决方案1】:

您可以使用窗口函数来实现这一点,假设列_1 是您的关键字,_2 是关键字的计数。在这种情况下 k = 2

scala> df.show()
+---+---+
| _1| _2|
+---+---+
|  1|  3|
|  2|  2|
|  1|  4|
|  1|  1|
|  2|  0|
|  1| 10|
|  2|  5|
+---+---+

scala> df.select('*,row_number().over(Window.orderBy('_2.desc).partitionBy('_1)).as("rn")).where('rn < 3).show()
+---+---+---+
| _1| _2| rn|
+---+---+---+
|  1| 10|  1|
|  1|  4|  2|
|  2|  5|  1|
|  2|  2|  2|
+---+---+---+

【讨论】:

  • 我正在考虑使用 Window 函数,但想知道 mapPartition 在计算方面是否会具有更好的性能?
猜你喜欢
  • 1970-01-01
  • 2021-05-10
  • 2013-03-15
  • 1970-01-01
  • 1970-01-01
  • 2013-10-30
  • 1970-01-01
  • 2016-05-13
  • 2019-05-01
相关资源
最近更新 更多