【发布时间】: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_VALUE 和 java.lang.ArrayIndexOutOfBoundException
我认为那是因为我的组太大了。
那么我该如何解决呢?有什么通用的方法可以解决这个问题吗?
【问题讨论】:
标签: java scala apache-spark indexoutofboundsexception illegalargumentexception