From f105c7d5f438f102d5f60a293a883c25ad078467 Mon Sep 17 00:00:00 2001 From: Andrey Yegorov Date: Fri, 23 Apr 2021 08:48:35 -0700 Subject: [PATCH 1/3] Fixed unnecessary copy to heap, see https://github.com/apache/pulsar/pull/10330 --- .../apache/bookkeeper/proto/checksum/DigestManager.java | 7 ++++++- .../main/java/org/apache/bookkeeper/util/ByteBufList.java | 7 ++++++- 2 files changed, 12 insertions(+), 2 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/checksum/DigestManager.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/checksum/DigestManager.java index 034dd6ed77e..92858554ef9 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/checksum/DigestManager.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/checksum/DigestManager.java @@ -19,6 +19,7 @@ import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBufAllocator; +import io.netty.buffer.CompositeByteBuf; import io.netty.buffer.Unpooled; import java.security.GeneralSecurityException; @@ -110,7 +111,11 @@ public ByteBufList computeDigestAndPackageForSending(long entryId, long lastAddC headersBuffer.writeLong(length); update(headersBuffer); - update(data); + if (data instanceof CompositeByteBuf) { + ((CompositeByteBuf) data).forEach(this::update); + } else { + update(data); + } populateValueAndReset(headersBuffer); return ByteBufList.get(headersBuffer, data); 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..21235d8e733 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 @@ -23,6 +23,7 @@ import com.google.common.annotations.VisibleForTesting; import io.netty.buffer.ByteBuf; +import io.netty.buffer.CompositeByteBuf; import io.netty.buffer.Unpooled; import io.netty.channel.ChannelHandler.Sharable; import io.netty.channel.ChannelHandlerContext; @@ -136,7 +137,11 @@ private static ByteBufList get() { * Append a {@link ByteBuf} at the end of this {@link ByteBufList}. */ public void add(ByteBuf buf) { - buffers.add(buf); + if (buf instanceof CompositeByteBuf) { + ((CompositeByteBuf) buf).forEach(buffers::add); + } else { + buffers.add(buf); + } } /** From c25668cc4d0fdae7745e44931dbb5a9c6ce8ef5a Mon Sep 17 00:00:00 2001 From: Andrey Yegorov Date: Fri, 23 Apr 2021 13:06:02 -0700 Subject: [PATCH 2/3] unwrap DuplicatedByteBuf and similar --- .../proto/checksum/DigestManager.java | 16 ++++++--- .../apache/bookkeeper/util/ByteBufList.java | 35 ++++++++++++++++--- 2 files changed, 42 insertions(+), 9 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/checksum/DigestManager.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/checksum/DigestManager.java index 92858554ef9..87d854105c2 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/checksum/DigestManager.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/checksum/DigestManager.java @@ -21,6 +21,7 @@ import io.netty.buffer.ByteBufAllocator; import io.netty.buffer.CompositeByteBuf; import io.netty.buffer.Unpooled; +import io.netty.util.ReferenceCountUtil; import java.security.GeneralSecurityException; import java.security.NoSuchAlgorithmException; @@ -111,14 +112,21 @@ public ByteBufList computeDigestAndPackageForSending(long entryId, long lastAddC headersBuffer.writeLong(length); update(headersBuffer); - if (data instanceof CompositeByteBuf) { - ((CompositeByteBuf) data).forEach(this::update); + + // don't unwrap slices + final ByteBuf unwrapped = data.unwrap() != null && data.unwrap() instanceof CompositeByteBuf + ? data.unwrap() : data; + ReferenceCountUtil.retain(unwrapped); + ReferenceCountUtil.release(data); + + if (unwrapped instanceof CompositeByteBuf) { + ((CompositeByteBuf) unwrapped).forEach(this::update); } else { - update(data); + update(unwrapped); } populateValueAndReset(headersBuffer); - return ByteBufList.get(headersBuffer, data); + return ByteBufList.get(headersBuffer, unwrapped); } /** 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 21235d8e733..d136ff0c58b 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 @@ -137,10 +137,19 @@ private static ByteBufList get() { * Append a {@link ByteBuf} at the end of this {@link ByteBufList}. */ public void add(ByteBuf buf) { - if (buf instanceof CompositeByteBuf) { - ((CompositeByteBuf) buf).forEach(buffers::add); + final ByteBuf unwrapped = buf.unwrap() != null && buf.unwrap() instanceof CompositeByteBuf + ? buf.unwrap() : buf; + ReferenceCountUtil.retain(unwrapped); + ReferenceCountUtil.release(buf); + + if (unwrapped instanceof CompositeByteBuf) { + ((CompositeByteBuf) unwrapped).forEach(b -> { + ReferenceCountUtil.retain(b); + buffers.add(b); + }); + ReferenceCountUtil.release(unwrapped); } else { - buffers.add(buf); + buffers.add(unwrapped); } } @@ -148,7 +157,23 @@ public void add(ByteBuf buf) { * Prepend a {@link ByteBuf} at the beginning of this {@link ByteBufList}. */ public void prepend(ByteBuf buf) { - buffers.add(0, buf); + // don't unwrap slices + final ByteBuf unwrapped = buf.unwrap() != null && buf.unwrap() instanceof CompositeByteBuf + ? buf.unwrap() : buf; + ReferenceCountUtil.retain(unwrapped); + ReferenceCountUtil.release(buf); + + if (unwrapped instanceof CompositeByteBuf) { + CompositeByteBuf composite = (CompositeByteBuf) unwrapped; + for (int i = composite.numComponents() - 1; i >= 0; i--) { + ByteBuf b = composite.component(i); + ReferenceCountUtil.retain(b); + buffers.add(0, b); + } + ReferenceCountUtil.release(unwrapped); + } else { + buffers.add(0, unwrapped); + } } /** @@ -264,7 +289,7 @@ public ByteBufList retain() { @Override protected void deallocate() { for (int i = 0; i < buffers.size(); i++) { - buffers.get(i).release(); + ReferenceCountUtil.release(buffers.get(i)); } buffers.clear(); From 0975eedac95b7934a825b7e962c7bd65a688f650 Mon Sep 17 00:00:00 2001 From: Andrey Yegorov Date: Fri, 23 Apr 2021 16:18:52 -0700 Subject: [PATCH 3/3] added test --- .../bookkeeper/util/ByteBufListTest.java | 34 +++++++++++++++++++ 1 file changed, 34 insertions(+) diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/ByteBufListTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/ByteBufListTest.java index 19c841be4c9..65f51e28482 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/ByteBufListTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/ByteBufListTest.java @@ -23,6 +23,7 @@ import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBufAllocator; +import io.netty.buffer.CompositeByteBuf; import io.netty.buffer.PooledByteBufAllocator; import io.netty.buffer.Unpooled; import io.netty.channel.Channel; @@ -88,6 +89,39 @@ public void testDouble() throws Exception { assertEquals(b2.refCnt(), 0); } + @Test + public void testComposite() throws Exception { + ByteBuf b1 = PooledByteBufAllocator.DEFAULT.heapBuffer(128, 128); + b1.writerIndex(b1.capacity()); + ByteBuf b2 = PooledByteBufAllocator.DEFAULT.heapBuffer(128, 128); + b2.writerIndex(b2.capacity()); + + CompositeByteBuf composite = PooledByteBufAllocator.DEFAULT.compositeBuffer(); + composite.addComponent(b1); + composite.addComponent(b2); + + ByteBufList buf = ByteBufList.get(composite); + + // composite is unwrapped into two parts + assertEquals(2, buf.size()); + // and released + assertEquals(composite.refCnt(), 0); + + assertEquals(256, buf.readableBytes()); + assertEquals(b1, buf.getBuffer(0)); + assertEquals(b2, buf.getBuffer(1)); + + assertEquals(buf.refCnt(), 1); + assertEquals(b1.refCnt(), 1); + assertEquals(b2.refCnt(), 1); + + buf.release(); + + assertEquals(buf.refCnt(), 0); + assertEquals(b1.refCnt(), 0); + assertEquals(b2.refCnt(), 0); + } + @Test public void testClone() throws Exception { ByteBuf b1 = PooledByteBufAllocator.DEFAULT.heapBuffer(128, 128);