【问题标题】:NullPointerException in Spark RDD map when submitted as a spark job作为 Spark 作业提交时 Spark RDD 映射中的 NullPointerException
【发布时间】:2016-12-23 12:42:45
【问题描述】:

我们正在尝试提交 spark 作业(spark 2.0、hadoop 2.7.2),但由于某种原因,我们在 EMR 中收到了一个相当神秘的 NPE。一切都作为一个 scala 程序运行得很好,所以我们不确定是什么导致了这个问题。这是堆栈跟踪:

18:02:55,271 错误 Utils:91 - 中止任务 java.lang.NullPointerException 在 org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIterator.agg_doAggregateWithKeys$(未知来源) 在 org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIterator.processNext(未知来源) 在 org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43) 在 org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$8$$anon$1.hasNext(WholeStageCodegenExec.scala:370) 在 scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:438) 在 org.apache.spark.sql.execution.datasources.DefaultWriterContainer$$anonfun$writeRows$1.apply$mcV$sp(WriterContainer.scala:253) 在 org.apache.spark.sql.execution.datasources.DefaultWriterContainer$$anonfun$writeRows$1.apply(WriterContainer.scala:252) 在 org.apache.spark.sql.execution.datasources.DefaultWriterContainer$$anonfun$writeRows$1.apply(WriterContainer.scala:252) 在 org.apache.spark.util.Utils$.tryWithSafeFinallyAndFailureCallbacks(Utils.scala:1325) 在 org.apache.spark.sql.execution.datasources.DefaultWriterContainer.writeRows(WriterContainer.scala:258) 在 org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand$$anonfun$run$1$$anonfun$apply$mcV$sp$1.apply(InsertIntoHadoopFsRelationCommand.scala:143) 在 org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand$$anonfun$run$1$$anonfun$apply$mcV$sp$1.apply(InsertIntoHadoopFsRelationCommand.scala:143) 在 org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:70) 在 org.apache.spark.scheduler.Task.run(Task.scala:85) 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:27​​4) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617) 在 java.lang.Thread.run(Thread.java:745)

据我们所知,这是通过以下方法发生的:

def process(dataFrame: DataFrame, S3bucket: String) = {
  dataFrame.map(row =>
      "text|label"
  ).coalesce(1).write.mode(SaveMode.Overwrite).text(S3bucket)
}

我们已将其范围缩小到 map 函数,因为它在作为 spark 作业提交时有效:

def process(dataFrame: DataFrame, S3bucket: String) = {
  dataFrame.coalesce(1).write.mode(SaveMode.Overwrite).text(S3bucket)
}

有谁知道是什么导致了这个问题?另外,我们该如何解决呢?我们很困惑。

【问题讨论】:

  • 你试过没有coalesce()吗?
  • @gsamaras 不!但它似乎在没有合并的情况下工作。这是怎么回事?

标签: scala hadoop apache-spark distributed-computing bigdata


【解决方案1】:

我认为当工作人员尝试访问仅存在于驱动程序而非工作人员上的 SparkContext 对象时,您会得到一个由工作人员抛出的 NullPointerException

coalesce() 重新分区您的数据。当您只请求一个分区时,它会尝试将所有数据压缩到一个分区*。这可能会给您的应用程序的内存占用带来很大压力。

一般来说,最好不要只将分区缩小 1。

更多信息,请阅读:Spark NullPointerException with saveAsTextFilethis


【讨论】:

  • 我们使用 coalesce(1) 的原因是将所有数据写入单个文件而不是多个文件。有没有其他方法可以做到这一点?
  • @cscan 没有。也许增加您的内存设置可以让您的应用程序使用 1 个分区,但我发布的错误并没有表明类似的情况。您是否有理由希望它们位于 1 个文件中?
  • 这个错误发生在我们只用五条记录进行测试时——我认为它与内存使用无关。
  • 我也是@cscan!我也更新了我的答案。顺便说一句,如果您真的只想拥有一个文件,那么我的建议是运行该作业,将文件拆分为多个部分,然后将这些部分合并到一个文件中。但总的来说,只有一个文件意味着没有很多数据,所以也许根本就不需要 Spark..
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-09-25
  • 1970-01-01
  • 2017-06-25
  • 1970-01-01
  • 2016-11-09
  • 1970-01-01
相关资源
最近更新 更多