【问题标题】:Out of Memory Error when Reading large file in Spark 2.1.0在 Spark 2.1.0 中读取大文件时出现内存不足错误
【发布时间】:2017-10-03 10:24:27
【问题描述】:

我想使用 spark 将大 (51GB) XML 文件(在外部硬盘上)读入数据帧(使用 spark-xml plugin),进行简单的映射/过滤,重新排序,然后将其写回磁盘,如CSV 文件。

但无论我如何调整,我总是得到java.lang.OutOfMemoryError: Java heap space

我想了解为什么增加分区数并不能停止OOM错误

不应该把任务分成更多的部分,让每个单独的部分更小,不会导致内存问题吗?

(Spark 不可能试图将所有内容都塞入内存并在不合适时崩溃,对吧??)

我尝试过的事情:

  • 在读取和写入时重新分区/合并到(5,000 和 10,000 个分区)数据帧(初始值为 1,604)
  • 使用较少数量的执行程序(6、4,即使使用 2 个执行程序,我也会收到 OOM 错误!)
  • 减小分割文件的大小(默认为 33MB)
  • 提供大量 RAM(我拥有的全部)
  • spark.memory.fraction 增加到 0.8(默认为 0.6)
  • spark.memory.storageFraction 减少到 0.2(默认为 0.5)
  • spark.default.parallelism 设置为 30 和 40(我默认为 8)
  • spark.files.maxPartitionBytes设置为64M(默认为128M)

我所有的代码都在这里(注意我没有缓存任何东西):

val df: DataFrame = spark.sqlContext.read
  .option("mode", "DROPMALFORMED")
  .format("com.databricks.spark.xml")
  .schema(customSchema) // defined previously
  .option("rowTag", "row")
  .load(s"$pathToInputXML")

println(s"\n\nNUM PARTITIONS: ${df.rdd.getNumPartitions}\n\n")
// prints 1604

// i pass `numPartitions` as cli arguments
val df2 = df.coalesce(numPartitions)

// filter and select only the cols i'm interested in
val dsout = df2
  .where( df2.col("_TypeId") === "1" )
  .select(
    df("_Id").as("id"),
    df("_Title").as("title"),
    df("_Body").as("body"),
  ).as[Post]

// regexes to clean the text
val tagPat = "<[^>]+>".r
val angularBracketsPat = "><|>|<"
val whitespacePat = """\s+""".r


// more mapping
dsout
 .map{
  case Post(id,title,body,tags) =>

    val body1 = tagPat.replaceAllIn(body,"")
    val body2 = whitespacePat.replaceAllIn(body1," ")

    Post(id,title.toLowerCase,body2.toLowerCase, tags.split(angularBracketsPat).mkString(","))

}
.orderBy(rand(SEED)) // random sort
.write // write it back to disk
.option("quoteAll", true)
.mode(SaveMode.Overwrite)
.csv(output)

注意事项

  • 输入拆分非常小(仅 33MB),那么为什么我不能有 8 个线程,每个线程处理一个拆分?它真的不应该破坏我的记忆(我已经知道了

更新我编写了一个较短版本的代码,它只读取文件,然后读取forEachPartition(println)。

我得到同样的 OOM 错误:

val df: DataFrame = spark.sqlContext.read
  .option("mode", "DROPMALFORMED")
  .format("com.databricks.spark.xml")
  .schema(customSchema)
  .option("rowTag", "row")
  .load(s"$pathToInputXML")
  .repartition(numPartitions)

println(s"\n\nNUM PARTITIONS: ${df.rdd.getNumPartitions}\n\n")

df
  .where(df.col("_PostTypeId") === "1")
  .select(
   df("_Id").as("id"),
   df("_Title").as("title"),
   df("_Body").as("body"),
   df("_Tags").as("tags")
  ).as[Post]
  .map {
    case Post(id, title, body, tags) =>
      Post(id, title.toLowerCase, body.toLowerCase, tags.toLowerCase))
  }
  .foreachPartition { rdd =>
    if (rdd.nonEmpty) {
      println(s"HI! I'm an RDD and I have ${rdd.size} elements!")
    }
  }

P.S.:我使用的是 spark v 2.1.0。我的机器有 8 个内核和 16 GB 内存。

【问题讨论】:

  • 您检查过 Spark UI 中创建的分区的大小吗?
  • @Khozzy 这是我运行应用程序时得到的结果,其中 1604 个分区用于读取 DF,50 个分区用于写入 DF:screenshot-spark-ui
  • 是的,但在作业执行期间查看 UI。你会发现每个任务执行了多长时间以及你的 CPU 是如何被利用的(可能有落后者)。
  • 我无法完成单个任务,请查看 OOM 错误之前的 UI:screenshot.. 另外,我注意到磁盘 IO 在崩溃之前出现了很大的峰值.
  • 其实我确实完成了8个初始任务中的一些,但是当第9个任务开始时,前面组中的一些任务失败了。

标签: xml scala apache-spark apache-spark-2.0 apache-spark-xml


【解决方案1】:

因为您要存储 RDD 两次并且 您的逻辑必须像这样更改或使用 SparkSql 过滤

 val df: DataFrame = SparkFactory.spark.read
      .option("mode", "DROPMALFORMED")
      .format("com.databricks.spark.xml")
      .schema(customSchema) // defined previously
      .option("rowTag", "row")
      .load(s"$pathToInputXML")
      .coalesce(numPartitions)

    println(s"\n\nNUM PARTITIONS: ${df.rdd.getNumPartitions}\n\n")
    // prints 1604


    // regexes to clean the text
    val tagPat = "<[^>]+>".r
    val angularBracketsPat = "><|>|<"
    val whitespacePat = """\s+""".r

    // filter and select only the cols i'm interested in
     df
      .where( df.col("_TypeId") === "1" )
      .select(
        df("_Id").as("id"),
        df("_Title").as("title"),
        df("_Body").as("body"),
      ).as[Post]
      .map{
        case Post(id,title,body,tags) =>

          val body1 = tagPat.replaceAllIn(body,"")
          val body2 = whitespacePat.replaceAllIn(body1," ")

          Post(id,title.toLowerCase,body2.toLowerCase, tags.split(angularBracketsPat).mkString(","))

      }
      .orderBy(rand(SEED)) // random sort
      .write // write it back to disk
      .option("quoteAll", true)
      .mode(SaveMode.Overwrite)
      .csv(output)

【讨论】:

  • 把它全部变成一个 DF 并没有真正的帮助。我还有java.lang.OutOfMemoryError: Java heap space
【解决方案2】:

您可以通过在环境变量中添加以下内容来更改堆大小:

  1. 环境变量名称:_JAVA_OPTIONS
  2. 环境变量值:-Xmx512M -Xms512M

【讨论】:

    【解决方案3】:

    我在运行 spark-shell 时遇到了这个错误,因此我将驱动程序内存增加到了一个很高的数字。然后我就可以加载 XML。

    spark-shell --driver-memory 6G
    

    来源:https://github.com/lintool/warcbase/issues/246#issuecomment-249272263

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-04-23
      • 2015-07-27
      • 2018-10-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-26
      相关资源
      最近更新 更多