Ume*_*cha 3 shuffle apache-spark rdd apache-spark-sql spark-dataframe
嗨我有Spark作业,它对ORC数据进行一些处理,并使用Spark 1.4.0中引入的DataFrameWriter save()API存储ORC数据.我有以下代码使用重型shuffle内存.如何优化以下代码?它有什么问题吗?它正如预期的那样工作正常,因为GC停顿并且随机播放大量数据而导致内存问题,从而导致速度变慢.请指导我是Spark新手.提前致谢.
JavaRDD<Row> updatedDsqlRDD = orderedFrame.toJavaRDD().coalesce(1, false).map(new Function<Row, Row>() {
@Override
public Row call(Row row) throws Exception {
List<Object> rowAsList;
Row row1 = null;
if (row != null) {
rowAsList = iterate(JavaConversions.seqAsJavaList(row.toSeq()));
row1 = RowFactory.create(rowAsList.toArray());
}
return row1;
}
}).union(modifiedRDD);
DataFrame updatedDataFrame = hiveContext.createDataFrame(updatedDsqlRDD,renamedSourceFrame.schema());
updatedDataFrame.write().mode(SaveMode.Append).format("orc").partitionBy("entity", "date").save("baseTable");
Run Code Online (Sandbox Code Playgroud)
编辑:根据建议我尝试将上面的代码转换为以下使用mapPartitionsWithIndex()但我仍然看到数据改组它比上面的代码更好,但它仍然失败通过命中GC限制并抛出OOM或进入GC暂停很长时间和超时和YARN会杀死遗嘱执行人.我使用spark.storage.memoryFraction为0.5和spark.shuffle.memoryFraction为0.4我尝试使用默认值并更改了许多组合没有任何帮助请指导
JavaRDD<Row> indexedRdd = sourceRdd.cache().mapPartitionsWithIndex(new Function2<Integer, Iterator<Row>, Iterator<Row>>() {
@Override
public Iterator<Row> call(Integer ind, Iterator<Row> rowIterator) throws Exception {
List<Row> rowList = new ArrayList<>();
while (rowIterator.hasNext()) {
Row row = rowIterator.next();
List<Object> rowAsList = iterate(JavaConversions.seqAsJavaList(row.toSeq()));
Row updatedRow = RowFactory.create(rowAsList.toArray());
rowList.add(updatedRow);
}
return rowList.iterator();
}
}, true).coalesce(200,true);
Run Code Online (Sandbox Code Playgroud)
将RDD或Dataframe合并到单个分区意味着您的所有处理都在一台计算机上进行.出于各种原因,这不是一件好事:所有数据都必须在网络中进行混洗,没有更多的并行性等等.相反,你应该看看其他运算符,如reduceByKey,mapPartitions,或者除此之外还有其他什么将数据合并到一台机器上.
注意:看你的代码我不明白你为什么把它带到一台机器上,你可能只是删除那部分.
| 归档时间: |
|
| 查看次数: |
2244 次 |
| 最近记录: |