【问题标题】:Java Spark broadcast and join two RDDsJava Spark 广播并加入两个 RDD
【发布时间】:2019-01-01 15:06:02
【问题描述】:

我有一张大桌子JavaPairRDD<String, MySchema> RDD1,还有一张较小的JavaPairRDD<String, Double> RDD2。我想加入这两个RDD,我知道最好的方法是使RDD2成为广播变量,然后加入以减少洗牌。如何处理广播部分?我的意思是在广播之后,我会得到一个变量(A List,或 Set),它不再是一个 RDD。如何使用 RDD 加入广播变量?

// I ignored the parsing part, just simplified it as loading from the files. 
JavaPairRDD<String, MySchema> RDD1 = sc.textFile ("path_to_small_dataset");
JavaPairRDD<String, Double> RDD2 = sc.textFile("path_to_large_dataset"); 

// Broadcast RDD2
Set<Tuple2<String, Double>> set2 = new HashSet<>();
set2.addAll(RDD2.collect());

// now I have set2 and RDD1, how can I join them? 

【问题讨论】:

    标签: apache-spark join rdd broadcast


    【解决方案1】:

    假设你有两个要加入的 RDD,第一个小到可以放入每个 worker 的内存中(smallRDD),第二个不需要洗牌完全没有(largeRDD)。

    在加入之前,你必须确保将大 RDD[T] 转换为 RDD[(key, T)]。键表示连接操作期间使用的列。

    这段代码在 Scala 中应该可以解决问题(但基本原理在 Java 中是相同的)

    val smallLookup = sc.broadcast(smallRDD.collect.toMap)
    largeRDD.flatMap { case(key, value) =>
      smallLookup.value.get(key).map { otherValue =>
      (key, (value, otherValue))
     }
    }
    

    希望对你有帮助

    【讨论】:

      猜你喜欢
      • 2016-09-07
      • 2016-01-24
      • 1970-01-01
      • 2015-06-22
      • 1970-01-01
      • 2017-06-27
      • 2018-11-24
      • 2015-10-18
      • 1970-01-01
      相关资源
      最近更新 更多