【发布时间】: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