【问题标题】:How to efficiently join an arbitrary number of RDDs?如何有效地加入任意数量的 RDD?
【发布时间】:2017-09-14 14:40:51
【问题描述】:

使用RDD1.join(RDD2) 连接两个RDD 很简单。但是,如果我在 List<JavaRDD> 中保留任意数量的 RDD,我怎样才能有效地加入它们?

【问题讨论】:

    标签: java join apache-spark java-8


    【解决方案1】:

    首先请注意,您不能加入JavaRDD。您需要使用以下方法获取JavaPairRDD

    • groupBy()(或keyBy()
    • cartesian()
    • [flat]mapToPair()
    • zipWithIndex()(很有用,因为它会在没有索引的地方添加索引)

    然后,一旦你有了你的名单,你就可以像这样加入他们:

    JavaPairRDD<Integer, String> linesA = sc.parallelizePairs(Arrays.asList(
                                                new Tuple2<>(1, "a1"),
                                                new Tuple2<>(2, "a2"),
                                                new Tuple2<>(3, "a3"),
                                                new Tuple2<>(4, "a4")));
    JavaPairRDD<Integer, String> linesB = sc.parallelizePairs(Arrays.asList(
                                                new Tuple2<>(1, "b1"),
                                                new Tuple2<>(5, "b5"),
                                                new Tuple2<>(3, "b3")));
    JavaPairRDD<Integer, String> linesC = sc.parallelizePairs(Arrays.asList(
                                                new Tuple2<>(1, "c1"),
                                                new Tuple2<>(5, "c6"),
                                                new Tuple2<>(6, "c3")));
    
    // the list of RDDs
    List<JavaPairRDD<Integer, String>> allLines = Arrays.asList(linesA, linesB, linesC);
    
    // since we probably don't want to modify any of the datasets in the list, we will
    // copy the first one in a separate variable to keep the result
    JavaPairRDD<Integer, String> res = allLines.get(0);
    for (int i = 1; i < allLines.size(); ++i) {  // note we skip position 0 !
        res = res.join(allLines.get(i))
        /*[1]*/  .mapValues(tuple -> tuple._1 + ':' + tuple._2);
    }
    

    带有[1] 的那一行是重要的,因为它映射了一个

    JavaPairRDD&lt;Integer, Tuple2&lt;String,String&gt;&gt;回一个

    JavaPairRdd&lt;Integer,String&gt; 使其与进一步的连接兼容。

    根据克里斯夫的回答,这可以像这样放在“一行”中:

    JavaPairRDD<Integer, String> res;
    res = allLines.stream()
                  .reduce((rdd1, rdd2) -> rdd1.join(rdd2).mapValues(tup -> tup._1 + ':' + tup._2))
                  .get();  // get value from Optional<JavaPairRDD>
    

    最后,关于性能的一些想法。在上面的示例中,我使用字符串连接将连接的结果减少回相同类型的 RDD。如果您有很多 RDD,您可能可以通过使用带有 JavaPairRDD&lt;Integer, StringBuilder&gt; resfor loop 版本来加快速度,您可以在其中手动进行第一次连接。如果需要,我会发布更多详细信息。

    【讨论】:

      【解决方案2】:

      我不熟悉 JavaRDD 类/接口,但也许您可以使用 Java 8 中的高阶函数 reduce 解决这个问题,请参阅 https://docs.oracle.com/javase/tutorial/collections/streams/reduction.html

      final List<JavaRDD> list = getList(); // where getList is your list implementation containing JavaRDD instances
      
      // The JavaRDD class provides rdd() to get the RDD
      final JavaRDD rdd = list.stream().map(JavaRDD::rdd).reduce(RDD::join);
      

      String 类的示例如下:-

      Stream.of("foo", "bar", "baz").reduce(String::concat);
      

      哪个产生

      美食吧

      【讨论】:

      • 我想我理解了这个想法,但 JavaRDD:: 没有产生连接。
      • 我刚刚看了一下javadoc,如果rdd()方法返回的是RDD,那么可以使用map高阶函数,见修改
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2010-09-17
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-03-22
      • 1970-01-01
      相关资源
      最近更新 更多