【发布时间】:2015-07-01 19:51:50
【问题描述】:
我的出发点是在 Scala 中使用 Apache Spark 的 RDD[(key,value)]。 RDD 包含大约 1500 万个元组。每个键大约有 50+-20 个值。
现在我想为每个键取一个值(不管是哪个值)。我目前的做法如下:
- 通过键对RDD进行HashPartition。 (没有明显的偏差)
- 按键对元组进行分组,得到 RDD[(key, array of values)]]
- 取每个值数组的第一个
基本上是这样的:
...
candidates
.groupByKey()
.map(c => (c._1, c._2.head)
...
分组是昂贵的部分。它仍然很快,因为没有网络洗牌并且候选者在内存中,但我可以做得更快吗?
我的想法是直接在分区上工作,但我不确定我能从 HashPartition 中得到什么。如果我取每个分区的第一个元组,我将获得每个键,但可能会根据分区数获得单个键的多个元组?还是我会错过钥匙?
谢谢!
【问题讨论】:
标签: scala apache-spark