【发布时间】:2017-04-12 10:19:55
【问题描述】:
您好,我正在尝试运行由于 StackoverflowError 而经常失败的长时间 sparkjob。该作业读取 parquetfile 并在 foreach 循环中创建一个 rdd。在做了一些研究之后,我认为为每个 rdd 创建一个检查点将帮助我解决我的内存问题。 (我尝试了不同的内存、开销内存、并行性、重新分区,并找到了最适合该作业的设置,但有时它仍然会失败,具体取决于我们集群上的负载。)
现在是我真正的问题。我正在尝试创建检查点,首先读取镶木地板创建一个 RDD,然后对其进行缓存,运行检查点函数,然后首先调用该操作以使检查点发生。在我指定的路径中没有创建检查点,YARN UI 表示该阶段被跳过。谁能帮我理解这个问题:)
ctx.getSparkContext().setCheckpointDir("/tmp/checkpoints");
public static void writeHhidToCouchbase(DataFrameContext ctx, List<String> filePathsStrings) {
filePathsStrings
.forEach(filePath -> {
JavaPairRDD<String, String> rdd =
UidHhidPerDay.getParquetFromPath(ctx, filePath);
rdd.cache();
rdd.checkpoint();
rdd.first();
rdd.foreachPartition(p -> {
CrumbsClient client = getClient();
p.forEachRemaining(uids -> {
Crumbs crumbs = client.getAsync(uids._1)
.timeout(10, TimeUnit.SECONDS)
.toBlocking()
.first();
String hHid = uids._2;
if (hHid != null) {
crumbs.getOrCreateSingletonCrumb(HouseholdCrumb.class).setHouseholdId(hHid);
client.putSync(crumbs);
}
});
client.shutdown();
});
});
}
检查点在第一次迭代中创建一次,但不再创建。 韩国
【问题讨论】:
-
不确定是否相关,但看起来对
rdd.checkpoint()的调用是多余的,rdd 没有要截断的父谱系。因此对内存使用没有帮助。 -
@ImDarrenG 谢谢!你能告诉我如何在代码中完成吗?
-
你可以删除对 checkpoint() 的调用,它对你没有帮助。
-
好的,但是还有其他方法可以确保作业成功完成吗?您能否解释一下您对 rdd 的含义没有要截断的父系谱系 :) 我对火花有点陌生。韩国
-
请您提供有关 stackoverflowerror 的更多详细信息?
标签: java apache-spark