【问题标题】:rdd.checkpoint is skipped in spark jobrdd.checkpoint 在 spark 作业中被跳过
【发布时间】: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


【解决方案1】:

我的错误是分区实际上是创建的。我上面提到的“第一个”分区是一个目录,里面有分区。由于像 8f987639-d5c8-46b8-a1e0-37081f9f8e00 这样的目录名称,我感到困惑。然而,查看@ImDarrenG 的血统评论给了我更多的见解。我从第一个缓存和检查点的 RDD 创建了一个新的重新分区的 RDD。这使得应用程序更加稳定,没有失败。

JavaPairRDD<String, String> rdd =
          UidHhidPerDay.getParquetFromPath(ctx, filePath);
      rdd.cache();
      rdd.checkpoint();
      rdd.first();
      JavaPairRDD<String, String> rddToCompute = rdd.repartition(72);
      rddToCompute.foreachPartition...

【讨论】:

  • 很高兴您发现了这个问题。我很想看看没有 rdd.cache(); 是否稳定。 rdd.checkpoint(); rdd.first();但如果它没有坏 - 不要修理它!
  • 谢谢,它不稳定我总是会得到 x nr 的失败任务 foreachpartition :),但现在我没有得到:)
猜你喜欢
  • 2021-05-13
  • 1970-01-01
  • 1970-01-01
  • 2022-12-18
  • 2022-11-28
  • 1970-01-01
  • 2014-10-08
  • 2020-07-27
  • 1970-01-01
相关资源
最近更新 更多