【发布时间】:2019-01-16 12:07:01
【问题描述】:
我的数据集大约有 2000 万行,需要大约 8 GB 的 RAM。我正在使用 2 个执行器运行我的工作,每个执行器 10 GB RAM,每个执行器 2 个内核。由于进一步的转换,数据应该一次被缓存。
我需要根据 4 个字段减少重复项(选择任何重复项)。两个选项:使用groupBy 和使用repartition 和mapPartitions。第二种方法允许您指定分区数,因此在某些情况下可以更快地执行,对吧?
你能解释一下哪个选项的性能更好吗?两个选项的 RAM 消耗是否相同?
使用groupBy
dataSet
.groupBy(col1, col2, col3, col4)
.agg(
last(col5),
...
last(col17)
);
使用repartition 和mapPartitions
dataSet.sqlContext().createDataFrame(
dataSet
.repartition(parallelism, seq(asList(col1, col2, col3, col4)))
.toJavaRDD()
.mapPartitions(DatasetOps::reduce),
SCHEMA
);
private static Iterator<Row> reduce(Iterator<Row> itr) {
Comparator<Row> comparator = (row1, row2) -> Comparator
.comparing((Row r) -> r.getAs(name(col1)))
.thenComparing((Row r) -> r.getAs(name(col2)))
.thenComparingInt((Row r) -> r.getAs(name(col3)))
.thenComparingInt((Row r) -> r.getAs(name(col4)))
.compare(row1, row2);
List<Row> list = StreamSupport
.stream(Spliterators.spliteratorUnknownSize(itr, Spliterator.ORDERED), false)
.collect(collectingAndThen(toCollection(() -> new TreeSet<>(comparator)), ArrayList::new));
return list.iterator();
}
【问题讨论】:
标签: apache-spark apache-spark-sql apache-spark-dataset