【发布时间】:2016-05-26 10:36:03
【问题描述】:
我有一个 RDD[(Int, List[Int])] 在每个分区中都有一个唯一的整数键。假设数据被分区为
分区 1 -> (1, List1), (2, List2)
分区 2 -> (1, List3), (2, List4)
所以当我想查找索引 1 的值时,我希望拥有
分区 1 -> (List1)
分区 2 -> (List3)
但是返回类型应该是 RDD[List(Int)] 而不是 Array[List(Int)] 这意味着我仍然希望集群上的分布式集合不收集到驱动程序。
目前我正在使用 filter( case { (k, v) => k == key} ).map(_._2) 但我知道这不会进行查找,而是按顺序搜索。
我知道有查找方法,但它返回并且 Array 而不是 RDD。 IndexedRDD 也是如此。
那么有没有办法在 Spark 中做到这一点?
【问题讨论】:
-
@zero323,据我所知,您的建议是; rrd.mapPartitions(it => it.toMap()(key).iterator)。我认为这可能是解决方案。唯一的开销是将迭代器转换为映射。在那之后,我相信这将是一个恒定时间的查找。我会试试这个并发布结果。
标签: scala apache-spark rdd