【问题标题】:Keeping an index to a list as a distributed collection - Spark将列表的索引作为分布式集合保存 - Spark
【发布时间】: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


【解决方案1】:

我认为你可以使用mapPartitionsWithIndex。您可以从 rdd 的分区器中获取感兴趣的分区,并且在您的 (Int,Iterator)=>Iterator 函数中,如果它不是感兴趣的分区,您只需返回一个空迭代器。为简单起见,如果它是正确的分区,我可能会返回传入的迭代器,然后只需使用 RDD.filter 来简化代码。

这里操作的迭代器可能是惰性的,所以实际上你的函数只是应用于分区号,而不是从迭代器中获取结果。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-09-10
    • 2014-06-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-15
    相关资源
    最近更新 更多