【问题标题】:Efficiently take one value for each key out of a RDD[(key,value)]有效地从 RDD[(key,value)] 中为每个键取一个值
【发布时间】:2015-07-01 19:51:50
【问题描述】:

我的出发点是在 Scala 中使用 Apache Spark 的 RDD[(key,value)]。 RDD 包含大约 1500 万个元组。每个键大约有 50+-20 个值。

现在我想为每个键取一个值(不管是哪个值)。我目前的做法如下:

  1. 通过键对RDD进行HashPartition。 (没有明显的偏差)
  2. 按键对元组进行分组,得到 RDD[(key, array of values)]]
  3. 取每个值数组的第一个

基本上是这样的:

...
candidates
.groupByKey()
.map(c => (c._1, c._2.head)
...

分组是昂贵的部分。它仍然很快,因为没有网络洗牌并且候选者在内存中,但我可以做得更快吗?

我的想法是直接在分区上工作,但我不确定我能从 HashPartition 中得到什么。如果我取每个分区的第一个元组,我将获得每个键,但可能会根据分区数获得单个键的多个元组?还是我会错过钥匙?

谢谢!

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    reduceByKey 使用返回第一个参数的函数怎么样?像这样:

    candidates.reduceByKey((x, _) => x)
    

    【讨论】:

    • 产生正确的结果。让我看看它在性能方面有何变化,有什么预测吗?
    • 一般来说,它会显着影响性能。对于您的限制(每个键 30-70 个值),它应该不多。
    • 这有点hacky,因为reduce应该与交换和关联函数一起使用。
    • @user52045 通过将(x, _) => x 更改为(x, y) => x max y 很容易使其可交换
    猜你喜欢
    • 1970-01-01
    • 2015-03-22
    • 2016-05-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-17
    • 2015-11-22
    相关资源
    最近更新 更多