From 0c294ede9020b2d71b7234e5a49b54a8a8c45c8d Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Sun, 5 Jun 2022 10:43:44 +0800 Subject: [PATCH 1/5] Support shrinkage in TripleLongPriorityQueue --- .../collections/TripleLongPriorityQueue.java | 46 +++++++++++++++++-- .../TripleLongPriorityQueueTest.java | 36 +++++++++++++++ 2 files changed, 79 insertions(+), 3 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java index 1d8d909beae79..f2bb4cb11934f 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java @@ -19,6 +19,7 @@ package org.apache.pulsar.common.util.collections; import static com.google.common.base.Preconditions.checkArgument; +import com.google.common.annotations.VisibleForTesting; import io.netty.buffer.ByteBuf; import io.netty.buffer.PooledByteBufAllocator; @@ -31,16 +32,28 @@ public class TripleLongPriorityQueue implements AutoCloseable { private static final int SIZE_OF_LONG = 8; private static final int DEFAULT_INITIAL_CAPACITY = 16; + private static final float DEFAULT_SHRINK_FACTOR = 0.5f; // Each item is composed of 3 longs private static final int ITEMS_COUNT = 3; private static final int TUPLE_SIZE = ITEMS_COUNT * SIZE_OF_LONG; + /** + * Reserve 10% of the capacity when shrinking to avoid frequent expansion and shrinkage. + */ + private static final float RESERVATION_FACTOR = 0.9f; + private final ByteBuf buffer; + private final int initialCapacity; + private int capacity; private int size; + /** + * When size < capacity * shrinkFactor, may trigger shrinking + */ + private final float shrinkFactor; /** * Create a new priority queue with default initial capacity. @@ -49,14 +62,21 @@ public TripleLongPriorityQueue() { this(DEFAULT_INITIAL_CAPACITY); } + public TripleLongPriorityQueue(int initialCapacity, float shrinkFactor) { + checkArgument(shrinkFactor > 0); + this.initialCapacity = initialCapacity; + this.capacity = initialCapacity; + this.buffer = PooledByteBufAllocator.DEFAULT.directBuffer(initialCapacity * ITEMS_COUNT * SIZE_OF_LONG); + this.size = 0; + this.shrinkFactor = shrinkFactor; + } + /** * Create a new priority queue with a given initial capacity. * @param initialCapacity */ public TripleLongPriorityQueue(int initialCapacity) { - capacity = initialCapacity; - buffer = PooledByteBufAllocator.DEFAULT.directBuffer(initialCapacity * ITEMS_COUNT * SIZE_OF_LONG); - size = 0; + this(initialCapacity, DEFAULT_SHRINK_FACTOR); } /** @@ -77,6 +97,8 @@ public void close() { public void add(long n1, long n2, long n3) { if (size == capacity) { increaseCapacity(); + } else { + shrinkCapacity(); } put(size, n1, n2, n3); @@ -122,6 +144,7 @@ public void pop() { swap(0, size - 1); size--; siftDown(0); + shrinkCapacity(); } /** @@ -144,6 +167,7 @@ public int size() { public void clear() { this.buffer.clear(); this.size = 0; + shrinkCapacity(); } private void increaseCapacity() { @@ -152,6 +176,17 @@ private void increaseCapacity() { buffer.capacity(this.capacity * TUPLE_SIZE); } + private void shrinkCapacity() { + if (capacity > initialCapacity && size < capacity * shrinkFactor) { + int decreasingSize = (int) (capacity * shrinkFactor * RESERVATION_FACTOR); + if (decreasingSize <= 0 || capacity - decreasingSize <= DEFAULT_INITIAL_CAPACITY) { + return; + } + this.capacity = capacity - decreasingSize; + buffer.capacity(this.capacity * TUPLE_SIZE); + } + } + private void siftUp(int idx) { while (idx > 0) { int parentIdx = (idx - 1) / 2; @@ -229,4 +264,9 @@ private void swap(int idx1, int idx2) { buffer.setLong(i2 + 1 * SIZE_OF_LONG, tmp2); buffer.setLong(i2 + 2 * SIZE_OF_LONG, tmp3); } + + @VisibleForTesting + ByteBuf getBuffer() { + return buffer; + } } diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueueTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueueTest.java index 4cb1027e0a9ff..c3e9a8b45aa0a 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueueTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueueTest.java @@ -135,4 +135,40 @@ public void testCompareWithSamePrefix() { pq.close(); } + + @Test + public void testShrink() throws Exception { + int initialCapacity = 20; + int tupleSize = 3 * 8; + TripleLongPriorityQueue pq = new TripleLongPriorityQueue(initialCapacity, 0.5f); + triggerScaleOut(initialCapacity, pq); + + assertEquals(pq.size(), initialCapacity); + assertEquals(pq.getBuffer().capacity(), initialCapacity * tupleSize); + + pq.add(0, 0, 0); + // Scale out to capacity * 2 + int scaleCapacity = initialCapacity * 2; + assertEquals(pq.getBuffer().capacity(), scaleCapacity * tupleSize); + // Trigger shrinking + for (int i = 0; i < initialCapacity / 2 + 1; i++) { + pq.pop(); + } + int capacity = scaleCapacity - (int)(scaleCapacity * 0.5f * 0.9f); + assertEquals(pq.getBuffer().capacity(), capacity * tupleSize); + // Scale out to capacity * 2 + triggerScaleOut(initialCapacity, pq); + scaleCapacity = capacity * 2; + assertEquals(pq.getBuffer().capacity(), scaleCapacity * tupleSize); + // Trigger shrinking + pq.clear(); + capacity = scaleCapacity - (int)(scaleCapacity * 0.5f * 0.9f); + assertEquals(pq.getBuffer().capacity(), capacity * tupleSize); + } + + private void triggerScaleOut(int initialCapacity, TripleLongPriorityQueue pq) { + for (long i = 0; i < initialCapacity; i++) { + pq.add(i, i, i); + } + } } From 8ba7ce3153202f709069c105377f846619a968fe Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Sun, 5 Jun 2022 11:09:56 +0800 Subject: [PATCH 2/5] Add unit test --- .../common/util/collections/TripleLongPriorityQueue.java | 8 ++++++-- .../util/collections/TripleLongPriorityQueueTest.java | 9 ++++----- 2 files changed, 10 insertions(+), 7 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java index f2bb4cb11934f..1f23ff55d1c2a 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java @@ -179,10 +179,14 @@ private void increaseCapacity() { private void shrinkCapacity() { if (capacity > initialCapacity && size < capacity * shrinkFactor) { int decreasingSize = (int) (capacity * shrinkFactor * RESERVATION_FACTOR); - if (decreasingSize <= 0 || capacity - decreasingSize <= DEFAULT_INITIAL_CAPACITY) { + if (decreasingSize <= 0) { return; } - this.capacity = capacity - decreasingSize; + if (capacity - decreasingSize <= initialCapacity) { + this.capacity = initialCapacity; + } else { + this.capacity = capacity - decreasingSize; + } buffer.capacity(this.capacity * TUPLE_SIZE); } } diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueueTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueueTest.java index c3e9a8b45aa0a..bd3aef86ad122 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueueTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueueTest.java @@ -141,13 +141,12 @@ public void testShrink() throws Exception { int initialCapacity = 20; int tupleSize = 3 * 8; TripleLongPriorityQueue pq = new TripleLongPriorityQueue(initialCapacity, 0.5f); - triggerScaleOut(initialCapacity, pq); - - assertEquals(pq.size(), initialCapacity); + pq.add(0, 0, 0); + assertEquals(pq.size(), 1); assertEquals(pq.getBuffer().capacity(), initialCapacity * tupleSize); - pq.add(0, 0, 0); // Scale out to capacity * 2 + triggerScaleOut(initialCapacity, pq); int scaleCapacity = initialCapacity * 2; assertEquals(pq.getBuffer().capacity(), scaleCapacity * tupleSize); // Trigger shrinking @@ -167,7 +166,7 @@ public void testShrink() throws Exception { } private void triggerScaleOut(int initialCapacity, TripleLongPriorityQueue pq) { - for (long i = 0; i < initialCapacity; i++) { + for (long i = 0; i < initialCapacity + 1; i++) { pq.add(i, i, i); } } From 6abd98189e2d298f615f8b8448973182af47cf3c Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Sun, 5 Jun 2022 11:11:21 +0800 Subject: [PATCH 3/5] Remove unused code --- .../pulsar/common/util/collections/TripleLongPriorityQueue.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java index 1f23ff55d1c2a..5b07c19228d30 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java @@ -97,8 +97,6 @@ public void close() { public void add(long n1, long n2, long n3) { if (size == capacity) { increaseCapacity(); - } else { - shrinkCapacity(); } put(size, n1, n2, n3); From f3e2ec4c66efe3eef3af8f6f19364bd468705d6d Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Sun, 5 Jun 2022 11:50:02 +0800 Subject: [PATCH 4/5] style --- .../pulsar/common/util/collections/TripleLongPriorityQueue.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java index 5b07c19228d30..63b91f8038ac5 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java @@ -51,7 +51,7 @@ public class TripleLongPriorityQueue implements AutoCloseable { private int capacity; private int size; /** - * When size < capacity * shrinkFactor, may trigger shrinking + * When size < capacity * shrinkFactor, may trigger shrinking. */ private final float shrinkFactor; From 1fe8d3772a3c943852f1fc82cd046d00be46539c Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Mon, 6 Jun 2022 23:39:36 +0800 Subject: [PATCH 5/5] Address comments --- .../collections/TripleLongPriorityQueue.java | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java index 63b91f8038ac5..487c2284cffdf 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/TripleLongPriorityQueue.java @@ -44,7 +44,7 @@ public class TripleLongPriorityQueue implements AutoCloseable { */ private static final float RESERVATION_FACTOR = 0.9f; - private final ByteBuf buffer; + private ByteBuf buffer; private final int initialCapacity; @@ -55,6 +55,8 @@ public class TripleLongPriorityQueue implements AutoCloseable { */ private final float shrinkFactor; + private float shrinkThreshold; + /** * Create a new priority queue with default initial capacity. */ @@ -66,7 +68,8 @@ public TripleLongPriorityQueue(int initialCapacity, float shrinkFactor) { checkArgument(shrinkFactor > 0); this.initialCapacity = initialCapacity; this.capacity = initialCapacity; - this.buffer = PooledByteBufAllocator.DEFAULT.directBuffer(initialCapacity * ITEMS_COUNT * SIZE_OF_LONG); + this.shrinkThreshold = this.capacity * shrinkFactor; + this.buffer = PooledByteBufAllocator.DEFAULT.directBuffer(initialCapacity * TUPLE_SIZE); this.size = 0; this.shrinkFactor = shrinkFactor; } @@ -171,11 +174,12 @@ public void clear() { private void increaseCapacity() { // For bigger sizes, increase by 50% this.capacity += (capacity <= 256 ? capacity : capacity / 2); + this.shrinkThreshold = this.capacity * shrinkFactor; buffer.capacity(this.capacity * TUPLE_SIZE); } private void shrinkCapacity() { - if (capacity > initialCapacity && size < capacity * shrinkFactor) { + if (capacity > initialCapacity && size < shrinkThreshold) { int decreasingSize = (int) (capacity * shrinkFactor * RESERVATION_FACTOR); if (decreasingSize <= 0) { return; @@ -185,7 +189,12 @@ private void shrinkCapacity() { } else { this.capacity = capacity - decreasingSize; } - buffer.capacity(this.capacity * TUPLE_SIZE); + this.shrinkThreshold = this.capacity * shrinkFactor; + + ByteBuf newBuffer = PooledByteBufAllocator.DEFAULT.directBuffer(this.capacity * TUPLE_SIZE); + buffer.getBytes(0, newBuffer, size * TUPLE_SIZE); + buffer.release(); + this.buffer = newBuffer; } }