From 9e5c1c091098568cf73448dfc597d508cb590f0c Mon Sep 17 00:00:00 2001 From: Sijie Guo Date: Fri, 6 Sep 2019 02:23:07 -0700 Subject: [PATCH 1/4] Avoid using directBuffer directly *Motivation* We introduced a smart buffer allocator since 4.9.0. It falls back to allocate buffer on heap when running out direct memory so that the bookie can run in a degraded mode rather than crashes. However if we use `directBuffer` directly, it bypasses the fallback mechanism. So the bookie crashes when running out of direct memory. *Modifications* Use `buffer` instead of `directBuffer` and let the allocator decide what is the best to use. --- .../src/main/java/org/apache/bookkeeper/bookie/Bookie.java | 2 +- .../java/org/apache/bookkeeper/bookie/BufferedChannel.java | 2 +- .../main/java/org/apache/bookkeeper/bookie/EntryLogger.java | 6 +++--- .../org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java | 4 ++-- .../apache/bookkeeper/bookie/storage/ldb/WriteCache.java | 4 ++-- .../main/java/org/apache/bookkeeper/util/ByteBufList.java | 2 +- 6 files changed, 10 insertions(+), 10 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Bookie.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Bookie.java index 6125970997d..002bc6476ff 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Bookie.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Bookie.java @@ -1308,7 +1308,7 @@ public void recoveryAddEntry(ByteBuf entry, WriteCallback cb, Object ctx, byte[] } private ByteBuf createExplicitLACEntry(long ledgerId, ByteBuf explicitLac) { - ByteBuf bb = allocator.directBuffer(8 + 8 + 4 + explicitLac.capacity()); + ByteBuf bb = allocator.buffer(8 + 8 + 4 + explicitLac.capacity()); bb.writeLong(ledgerId); bb.writeLong(METAENTRY_ID_LEDGER_EXPLICITLAC); bb.writeInt(explicitLac.capacity()); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BufferedChannel.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BufferedChannel.java index 31fb2035ea9..3735cbc553e 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BufferedChannel.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BufferedChannel.java @@ -83,7 +83,7 @@ public BufferedChannel(ByteBufAllocator allocator, FileChannel fc, int writeCapa this.writeCapacity = writeCapacity; this.position = new AtomicLong(fc.position()); this.writeBufferStartPosition.set(position.get()); - this.writeBuffer = allocator.directBuffer(writeCapacity); + this.writeBuffer = allocator.buffer(writeCapacity); this.unpersistedBytes = new AtomicLong(0); this.unpersistedBytesBound = unpersistedBytesBound; this.doRegularFlushes = unpersistedBytesBound > 0; diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/EntryLogger.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/EntryLogger.java index 6662d594d27..b50c7e1a699 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/EntryLogger.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/EntryLogger.java @@ -861,7 +861,7 @@ private Header getHeaderForLogId(long entryLogId) throws IOException { BufferedReadChannel bc = getChannelForLogId(entryLogId); // Allocate buffer to read (version, ledgersMapOffset, ledgerCount) - ByteBuf headers = allocator.directBuffer(LOGFILE_HEADER_SIZE); + ByteBuf headers = allocator.buffer(LOGFILE_HEADER_SIZE); try { bc.read(headers, 0); @@ -977,7 +977,7 @@ public void scanEntryLog(long entryLogId, EntryLogScanner scanner) throws IOExce long pos = LOGFILE_HEADER_SIZE; // Start with a reasonably sized buffer size - ByteBuf data = allocator.directBuffer(1024 * 1024); + ByteBuf data = allocator.buffer(1024 * 1024); try { @@ -1059,7 +1059,7 @@ EntryLogMetadata extractEntryLogMetadataFromIndex(long entryLogId) throws IOExce EntryLogMetadata meta = new EntryLogMetadata(entryLogId); final int maxMapSize = LEDGERS_MAP_HEADER_SIZE + LEDGERS_MAP_ENTRY_SIZE * LEDGERS_MAP_MAX_BATCH_SIZE; - ByteBuf ledgersMap = allocator.directBuffer(maxMapSize); + ByteBuf ledgersMap = allocator.buffer(maxMapSize); try { while (offset < bc.size()) { diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java index 986c741a010..976c4137e76 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java @@ -73,7 +73,7 @@ public ReadCache(ByteBufAllocator allocator, long maxCacheSize, int maxSegmentSi cacheIndexes = new ArrayList<>(); for (int i = 0; i < segmentsCount; i++) { - cacheSegments.add(Unpooled.directBuffer(segmentSize, segmentSize)); + cacheSegments.add(Unpooled.buffer(segmentSize, segmentSize)); cacheIndexes.add(new ConcurrentLongLongPairHashMap(4096, 2 * Runtime.getRuntime().availableProcessors())); } } @@ -142,7 +142,7 @@ public ByteBuf get(long ledgerId, long entryId) { int entryOffset = (int) res.first; int entryLen = (int) res.second; - ByteBuf entry = allocator.directBuffer(entryLen, entryLen); + ByteBuf entry = allocator.buffer(entryLen, entryLen); entry.writeBytes(cacheSegments.get(segmentIdx), entryOffset, entryLen); return entry; } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java index ac58e8eacac..eb5f6d9ff2c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java @@ -105,11 +105,11 @@ public WriteCache(ByteBufAllocator allocator, long maxCacheSize, int maxSegmentS for (int i = 0; i < segmentsCount - 1; i++) { // All intermediate segments will be full-size - cacheSegments[i] = Unpooled.directBuffer(maxSegmentSize, maxSegmentSize); + cacheSegments[i] = allocator.buffer(maxSegmentSize, maxSegmentSize); } int lastSegmentSize = (int) (maxCacheSize % maxSegmentSize); - cacheSegments[segmentsCount - 1] = Unpooled.directBuffer(lastSegmentSize, lastSegmentSize); + cacheSegments[segmentsCount - 1] = allocator.buffer(lastSegmentSize, lastSegmentSize); } public void clear() { diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/ByteBufList.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/ByteBufList.java index 355cf3f307b..d7d047efede 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/ByteBufList.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/ByteBufList.java @@ -306,7 +306,7 @@ public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) if (prependSize) { // Prepend the frame size before writing the buffer list, so that we only have 1 single size // header - ByteBuf sizeBuffer = ctx.alloc().directBuffer(4, 4); + ByteBuf sizeBuffer = ctx.alloc().buffer(4, 4); sizeBuffer.writeInt(b.readableBytes()); ctx.write(sizeBuffer, ctx.voidPromise()); } From 1615df3e3e6ff725a7325b36e953fbb5c81e7fea Mon Sep 17 00:00:00 2001 From: Sijie Guo Date: Mon, 9 Sep 2019 12:35:57 -0700 Subject: [PATCH 2/4] Fix checkstyle issue --- .../org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java | 1 - 1 file changed, 1 deletion(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java index eb5f6d9ff2c..1d5365081c4 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/WriteCache.java @@ -24,7 +24,6 @@ import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBufAllocator; -import io.netty.buffer.Unpooled; import java.io.Closeable; import java.util.concurrent.atomic.AtomicLong; From 93b16abfc412e7cb053f88268f0cbb8519c8c969 Mon Sep 17 00:00:00 2001 From: Sijie Guo Date: Tue, 10 Sep 2019 16:31:18 -0700 Subject: [PATCH 3/4] Use allocator consistently --- .../org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java index 976c4137e76..77bcbe840f6 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java @@ -73,7 +73,7 @@ public ReadCache(ByteBufAllocator allocator, long maxCacheSize, int maxSegmentSi cacheIndexes = new ArrayList<>(); for (int i = 0; i < segmentsCount; i++) { - cacheSegments.add(Unpooled.buffer(segmentSize, segmentSize)); + cacheSegments.add(allocator.buffer(segmentSize, segmentSize)); cacheIndexes.add(new ConcurrentLongLongPairHashMap(4096, 2 * Runtime.getRuntime().availableProcessors())); } } From f5134ed4828801925de9d1379d560b8983f90030 Mon Sep 17 00:00:00 2001 From: Sijie Guo Date: Wed, 11 Sep 2019 18:21:18 -0700 Subject: [PATCH 4/4] Fix checkstyle --- .../java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java | 1 - 1 file changed, 1 deletion(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java index 77bcbe840f6..1f49e62a96b 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/ReadCache.java @@ -24,7 +24,6 @@ import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBufAllocator; -import io.netty.buffer.Unpooled; import java.io.Closeable; import java.util.ArrayList;