Reactor unzip files in parallel to memory leads to OutOfMemoryError

Viewed 256

I'm trying to process a huge tar file containing multiple inner XML files to be processed. I have disk constraints, so I have to use an uncompress-in-memory-do-stuff-then-forget approach. I'm using Apache Compress library. So far so good, this is working fine, but I would like now to parallelise the processing of the files to improve the execution time. To do so, we are using reactor.

This looks to be working until, for some reason, data is not been cleaned up after being processed, so it ends up in a java.lang.OutOfMemoryError: Java heap space exception. Other things I have considered is to tune the garbage collector to step in faster, but I'm a bit lost so I appreciate any kind of help.

Getting into the code. This is the Flux.generate I'm using. For every next requested item, it generates an array of bytes containing one uncompressed file from the tar file. Looks like this:

(ArchiveInputStream tarArchive, SynchronousSink<byte[]> sink) -> {
    process: try {
      ArchiveEntry entry = null;
      while ((entry = tarArchive.getNextEntry()) != null) {
        if (entry.isDirectory()) continue;
        log.info(entry.getName());

        var bytes = IOUtils.readRange(tarArchive, (int) entry.getSize());
        sink.next(bytes);

        break process;
      }
      log.info("done!!!!");
      sink.complete();
    } catch (IOException e) {
      log.error("Exception getting file entry", e);
      sink.error(e);
    }
    return tarArchive;
  }

Later on, the array of bytes is converted into a ByteArrayInputStream to be mapped by JAXB into the internal model:

  public Feed processElement(byte[] file) {
    var feed = new Feed();
    if (file.length != 0) {
      try(var is = new ByteArrayInputStream(file)) {
        var jaxbUnmarshaller = jaxbContext.createUnmarshaller();
        feed = (Feed) jaxbUnmarshaller.unmarshal(is);
      } catch (Exception ex) {
        log.error("Error during XML unmarshal", ex);
      }
    }
    return feed;
  }

And stream is parallelised as follows:

uncompressFiles(tarFileName)
        .parallel()
        .runOn(parallel())
        .map(parser::processElement)
        .sequential()
//.....

Let me know if you have any insight either on in this particular case or about how to resolve memory leaks.

1 Answers

After some days of research, it turns out that there is no memory leak at all. Issue was caused by the buffering nature of reactive streams. So after adding some traces I discovered that could be hundreds of objects in flight between a file is read and sinked. This is a reasonable performance optimisation for most of the cases but in this particular one objects are too heavy (100 files could weight easily 8GB)

What I did in first place was to reduce the buffer size with the System.property reactor.bufferSize.small to 16 (minimum possible).

Though this is not enough because of the parallel step. So every parallel thread ask to the Producer 16 files * 12 threads (on my macbook) = 192 files on flight. To avoid this we can limit the fetch rate of every thread as follows

.runOn(parallel(), 16 / DEFAULT_POOL_SIZE)

So we ensure that no more than 16 files are on memory at a time for this step. This way I could finally keep the memory peeks under 4GB without loosing any processing power.

Related