From 8b3bf4de44a89d9ace716ccbe49a086aadbc6606 Mon Sep 17 00:00:00 2001 From: "Maxim.Mossienko" Date: Wed, 14 Mar 2012 17:18:04 +0400 Subject: [PATCH] limit pending write requests size to avoid OOME --- .../util/io/storage/RefCountingStorage.java | 47 ++++++++++++------- 1 file changed, 29 insertions(+), 18 deletions(-) diff --git a/platform/util/src/com/intellij/util/io/storage/RefCountingStorage.java b/platform/util/src/com/intellij/util/io/storage/RefCountingStorage.java index eb8ba08a3140..a1e015453c9b 100644 --- a/platform/util/src/com/intellij/util/io/storage/RefCountingStorage.java +++ b/platform/util/src/com/intellij/util/io/storage/RefCountingStorage.java @@ -36,6 +36,7 @@ import java.util.zip.InflaterInputStream; public class RefCountingStorage extends AbstractStorage { private final Map> myPendingWriteRequests = new ConcurrentHashMap>(); + private int myPendingWriteRequestsSize; private final ThreadPoolExecutor myPendingWriteRequestsExecutor = new ThreadPoolExecutor(1, 1, Long.MAX_VALUE, TimeUnit.DAYS, new LinkedBlockingQueue(), new ThreadFactory() { @Override public Thread newThread(Runnable runnable) { @@ -44,6 +45,7 @@ public class RefCountingStorage extends AbstractStorage { }); private final boolean myDoNotZipCaches = Boolean.valueOf(System.getProperty("idea.doNotZipCaches")).booleanValue(); + private static final int MAX_PENDING_WRITE_SIZE = 20 * 1024 * 1024; public RefCountingStorage(String path) throws IOException { super(path); @@ -96,26 +98,35 @@ public class RefCountingStorage extends AbstractStorage { waitForPendingWriteForRecord(record); synchronized (myLock) { - Future future = myPendingWriteRequestsExecutor.submit(new Callable() { - @Override - public Object call() throws IOException { - BufferExposingByteArrayOutputStream s = new BufferExposingByteArrayOutputStream(); - DeflaterOutputStream out = new DeflaterOutputStream(s); - try { - out.write(bytes.getBytes(), bytes.getOffset(), bytes.getLength()); - } - finally { - out.close(); + myPendingWriteRequestsSize += bytes.getLength(); + if (myPendingWriteRequestsSize > MAX_PENDING_WRITE_SIZE) { + zipAndWrite(bytes, record, fixedSize); + } else { + myPendingWriteRequests.put(record, myPendingWriteRequestsExecutor.submit(new Callable() { + @Override + public Object call() throws IOException { + zipAndWrite(bytes, record, fixedSize); + return null; } + })); + } + } + } - synchronized (myLock) { - doWrite(record, fixedSize, s); - myPendingWriteRequests.remove(record); - } - return null; - } - }); - myPendingWriteRequests.put(record, future); + private void zipAndWrite(ByteSequence bytes, int record, boolean fixedSize) throws IOException { + BufferExposingByteArrayOutputStream s = new BufferExposingByteArrayOutputStream(); + DeflaterOutputStream out = new DeflaterOutputStream(s); + try { + out.write(bytes.getBytes(), bytes.getOffset(), bytes.getLength()); + } + finally { + out.close(); + } + + synchronized (myLock) { + doWrite(record, fixedSize, s); + myPendingWriteRequestsSize -= bytes.getLength(); + myPendingWriteRequests.remove(record); } }