Is there a better way to load a huge tar file in Spark while avoiding OutOfMemoryError?

Viewed 206

I have a single tar file mytar.tar that is 40 GB in size. Inside this tar file are 500 tar.gz files, and inside each one of these tar.gz files are a bunch of JSON files. I have written up the code to process this single tar file and attempt to get the list of JSON string contents. My code looks like the following.

val isRdd = sc.binaryFiles("/mnt/mytar.tar")
  .flatMap(t => { 
    val buf = scala.collection.mutable.ListBuffer.empty[TarArchiveInputStream]
    val stream = t._2
    val is = new TarArchiveInputStream(stream.open())
    var entry = is.getNextTarEntry()
    while (entry != null) {
      val name = entry.getName()
      val size = entry.getSize.toInt

      if (entry.isFile() && size > -1) {
        val content = new Array[Byte](size)
        is.read(content, 0, content.length)

        val tgIs = new TarArchiveInputStream(new GzipCompressorInputStream(new ByteArrayInputStream(content)))
        buf += tgIs
      }
      entry = is.getNextTarEntry()
    }
    buf.toList
  })
  .cache

val byteRdd = isRdd.flatMap(is => {
    val buf = scala.collection.mutable.ListBuffer.empty[Array[Byte]]
    var entry = is.getNextTarEntry()
    while (entry != null) {
      val name = entry.getName()
      val size = entry.getSize.toInt

      if (entry.isFile() && name.endsWith(".json") && size > -1) {
        val data = new Array[Byte](size)
        is.read(data, 0, data.length)
        buf += data
      }
      entry = is.getNextTarEntry()
    }
    buf.toList
  })
  .cache

val jsonRdd = byteRdd
  .map(arr => getJson(arr))
  .filter(_.length > 0)
  .cache

jsonRdd.count //action just to execute the code

When I execute this code, I get an OutOfMemoryError (OOME).

org.apache.spark.SparkException: Job aborted due to stage failure: 
Task 0 in stage 24.0 failed 4 times, most recent failure: 
Lost task 0.3 in stage 24.0 (TID 137, 10.162.224.171, executor 13): 
java.lang.OutOfMemoryError: Java heap space

My EC2 cluster has 1 driver and 2 worker nodes of type i3.xlarge (30.5 GB Memory, 4 Cores). From looking at the logs and just thinking about it, I believe the OOME happens during the creation of the isRDD (input stream RDD).

Is there anything else in the code or creation of my Spark cluster that I can do to mitigate this problem? Should I select an EC2 instance with more memory (e.g. a memory optimized instance like R5.2xlarge)? FWIW, I upgraded to an R5.2xlarge cluster setting and still saw the OOME.

One thing I have thought about doing was to untar mytar.tar and instead start with the .tar.gz files inside. I am thinking that each .tar.gz inside the tar file will have to be less than 30 GB to avoid the OOME (on the i3.xlarge).

Any tips or advice is appreciated.

0 Answers
Related