【问题标题】:Spark: Cogroup RDDs fails in case of huge groupSpark:在大组的情况下,Cogroup RDD 失败
【发布时间】:2015-06-18 15:03:37
【问题描述】:

下午好!我有一个问题:

val rdd1: RDD[(key, value)] = ...
val rdd2: RDD[(key, othervalue)] = ...

我想过滤rdd1 并丢弃所有不在rdd2 中的元素。我知道有两种方法可以做到这一点。

第一:

val keySet = rdd2.map(_.key).distinct().collect().toSet
rdd1.filter(x => keySet.contains(x))

它不起作用,因为keySet 很大并且不适合内存。

另一个:

rdd1.cogroup(rdd2)
  .filter(x => x._2._2.nonEmpty)
  .flatMap(x => x._2._1)

这里发生了一些事情,我得到了两种错误(在不同的代码位置):java.lang.IllegalArgumentException: Size exceeds Integer.MAX_VALUEjava.lang.ArrayIndexOutOfBoundException

我认为那是因为我的组太大了。

那么我该如何解决呢?有什么通用的方法可以解决这个问题吗?

【问题讨论】:

    标签: java scala apache-spark indexoutofboundsexception illegalargumentexception


    【解决方案1】:

    您考虑过使用subtractByKey 吗?

    类似

    的东西
    rdd1.map(x => (x, x))
        .subtractByKey(rdd2)
        .map((k,v) => k)
    

    【讨论】:

      【解决方案2】:

      考虑 rdd1.subtractByKey( rdd1.subtractByKey(rdd2) )。 rdd1.subtractByKey(rdd2) 将获取那些键在 rdd1 中但不在 rdd2 中的元素。这与你想要的相反。减去那些将得到你想要的。

      【讨论】:

        猜你喜欢
        • 2016-04-26
        • 2022-01-22
        • 2016-01-07
        • 1970-01-01
        • 2016-09-07
        • 2014-11-07
        • 1970-01-01
        • 2016-01-13
        • 1970-01-01
        相关资源
        最近更新 更多