【发布时间】:2017-09-14 14:40:51
【问题描述】:
使用RDD1.join(RDD2) 连接两个RDD 很简单。但是,如果我在 List<JavaRDD> 中保留任意数量的 RDD,我怎样才能有效地加入它们?
【问题讨论】:
标签: java join apache-spark java-8
使用RDD1.join(RDD2) 连接两个RDD 很简单。但是,如果我在 List<JavaRDD> 中保留任意数量的 RDD,我怎样才能有效地加入它们?
【问题讨论】:
标签: java join apache-spark java-8
首先请注意,您不能加入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<Integer, Tuple2<String,String>>回一个
JavaPairRdd<Integer,String> 使其与进一步的连接兼容。
根据克里斯夫的回答,这可以像这样放在“一行”中:
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<Integer, StringBuilder> res 的 for loop 版本来加快速度,您可以在其中手动进行第一次连接。如果需要,我会发布更多详细信息。
【讨论】:
我不熟悉 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);
哪个产生
美食吧
【讨论】:
rdd()方法返回的是RDD,那么可以使用map高阶函数,见修改