【问题标题】:Find RDD[(T, U)] elements that have key in RDD[T]查找在 RDD[T] 中有键的 RDD[(T, U)] 元素
【发布时间】:2016-03-18 23:48:34
【问题描述】:

给定

val as: RDD[(T, U)]
val bs: RDD[T]

我想过滤as 以查找存在bs 键的元素。

一种方法是

val intermediateAndOtherwiseUnnessaryPair = bs.map(b => b -> b)
bs.join(as).values

但是bs 上的映射很不幸。有没有更直接的方法?

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    您可以通过以下方式减少映射的不必要性:

    val intermediateAndOtherwiseUnnessaryPair = bs.map(b => (b, 1))
    

    此外,加入前的共同分区也有很大帮助:

    val intermediateAndOtherwiseUnnessaryPair = bs.map(b => (b, 1)).paritionBy(new HashPartitioner(NUM_PARTITIONS))
    
    bs.paritionBy(new HashPartitioner(NUM_PARTITIONS)).join(as).values
    

    Co 分区的 RDD 在运行时不会被洗牌,因此您会看到显着的性能提升。

    如果bs 太大(更准确地说,具有大量唯一值),广播可能不起作用,您可能还想增加driver.maxResultsize

    【讨论】:

    • 您对此有任何更改基准吗?我的意思是使用 partitionBy 与不使用它?在我看来,asbs 在创建时仍然会随机播放。
    • @MateuszDymczyk 没错,洗牌是不可避免的。但是,我在实践中看到进行预分区可以避免所有“FetchFailedExceptions”,我认为这是由内存问题引起的。在黑暗中拍摄,但感觉在 RDD 和 join 之间移动的工作量太大,预分区使框架更容易。同样,这一切都来自经验,没有基准。
    • FetchFailedException 可能由于许多原因而发生,例如执行程序没有响应或您自己(或默认情况下由 Spark)错误地设置了分区数,但这实际上取决于用例。如果你设置了太多的分区,你实际上会得到更多的随机播放。我想说这是过早的优化,在很多情况下可能会适得其反。
    • 我明白了。但是,我可以重现 join 在我处理的数据集(~1T)上没有先进行共同分区的情况下无法工作的场景。此外,有趣的一点是由于分区数不正确导致 fetchfailed。为什么会这样?根据我的经验,增加分区数量(~5 * 核心数)通常会有所帮助,而减少分区数量从来没有帮助过。但我猜想 spark 有太多的旋钮和句柄来做一般性的陈述。
    • @EndeNeu 不,如果 RDD 已经分区,则洗牌不会发生两次:safaribooksonline.com/library/view/learning-spark/9781449359034/…
    【解决方案2】:

    使用第二个RDD 过滤一个RDD 的唯一两种(或至少我知道的唯一一种)流行且通用的方法是:

    1) join 你已经在做 - 在这种情况下,我不会担心不必要的中间 RDD 不过,map() 是一个狭窄的转换,不会引入那么多开销。不过,join() 本身很可能会很慢,因为它是一个广泛的转换(需要洗牌)

    2) 收集驱动程序上的bs 并使其成为广播变量,然后将在as.filter() 中使用

    val collected = sc.broadcast(bs.collect().toSet)
    as.filter(el => collected.value.contains(el))
    

    您需要这样做,因为 Spark 不支持在 RDD 上调用的方法中嵌套 RDDs

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多