From 4a632cdba0f78afb8adf497eb96e8733764028c4 Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 1 Jul 2020 16:15:39 +0800 Subject: [PATCH 1/3] Fix batch ackset recycled multiple times. --- .../client/impl/BatchMessageIndexAckTest.java | 31 +++++++++++++++++++ .../util/collections/BitSetRecyclable.java | 8 +++-- 2 files changed, 37 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java index 114b9ae3fdc1e..f24678202b924 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java @@ -181,4 +181,35 @@ public void testBatchMessageIndexAckForExclusiveSubscription() throws PulsarClie // broker also need to handle the available permits. Assert.assertEquals(received.size(), 100); } + + @Test + public void testSafeAckSetRecycle() throws Exception { + final String topic = "persistent://my-property/my-ns/testSafeAckSetRecycle"; + + Producer producer = pulsarClient.newProducer() + .batchingMaxMessages(10) + .blockIfQueueFull(true).topic(topic) + .create(); + + Consumer consumer = pulsarClient.newConsumer() + .acknowledgmentGroupTime(1, TimeUnit.MILLISECONDS) + .topic(topic) + .subscriptionName("test") + .subscribe(); + + final int messages = 100; + for (int i = 0; i < messages; i++) { + producer.sendAsync("Hello Pulsar".getBytes()); + } + + // Should not throw an exception. + for (int i = 0; i < messages; i++) { + consumer.acknowledgeCumulative(consumer.receive()); + // make sure the group ack flushed. + Thread.sleep(2); + } + + producer.close(); + consumer.close(); + } } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java index 285416ac606b2..0c0031bb90325 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java @@ -27,6 +27,7 @@ import java.nio.LongBuffer; import java.util.Arrays; import java.util.BitSet; +import java.util.concurrent.atomic.AtomicBoolean; /** * This this copy of {@link BitSet}. @@ -1183,6 +1184,7 @@ public BitSetRecyclable resetWords(long[] words) { return this; } + private final AtomicBoolean recycled = new AtomicBoolean(false); private Handle recyclerHandle = null; private static final Recycler RECYCLER = new Recycler() { @@ -1197,11 +1199,13 @@ private BitSetRecyclable(Handle recyclerHandle) { } public static BitSetRecyclable create() { - return RECYCLER.get(); + BitSetRecyclable instance = RECYCLER.get(); + instance.recycled.set(false); + return instance; } public void recycle() { - if (recyclerHandle != null) { + if (recyclerHandle != null && recycled.compareAndSet(false, true)) { this.clear(); recyclerHandle.recycle(this); } From 067c4509dc9615c969557178180cb44bf47cc15b Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 2 Jul 2020 15:37:25 +0800 Subject: [PATCH 2/3] Apply comments. --- .../pulsar/client/impl/BatchMessageIndexAckTest.java | 2 +- .../impl/PersistentAcknowledgmentsGroupingTracker.java | 1 + .../java/org/apache/pulsar/common/protocol/Commands.java | 3 --- .../pulsar/common/util/collections/BitSetRecyclable.java | 7 ++----- 4 files changed, 4 insertions(+), 9 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java index f24678202b924..3150f10079ca7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java @@ -183,7 +183,7 @@ public void testBatchMessageIndexAckForExclusiveSubscription() throws PulsarClie } @Test - public void testSafeAckSetRecycle() throws Exception { + public void testDoNotRecycleAckSetMultipleTimes() throws Exception { final String topic = "persistent://my-property/my-ns/testSafeAckSetRecycle"; Producer producer = pulsarClient.newProducer() diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java index 937f005e2cff5..690897998b83d 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java @@ -205,6 +205,7 @@ private boolean doImmediateBatchIndexAck(BatchMessageIdImpl msgId, int batchInde } final ByteBuf cmd = Commands.newAck(consumer.consumerId, msgId.ledgerId, msgId.entryId, bitSet, ackType, null, properties); + bitSet.recycle(); cnx.ctx().writeAndFlush(cmd, cnx.ctx().voidPromise()); return true; } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index 0222f3c30b91f..7ae21e7131991 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -961,9 +961,6 @@ public static ByteBuf newAck(long consumerId, long ledgerId, long entryId, BitSe ByteBuf res = serializeWithSize(BaseCommand.newBuilder().setType(Type.ACK).setAck(ack)); ack.recycle(); - if (ackSet != null) { - ackSet.recycle(); - } ackBuilder.recycle(); messageIdDataBuilder.recycle(); messageIdData.recycle(); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java index 0c0031bb90325..c586a54d46f77 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java @@ -1184,7 +1184,6 @@ public BitSetRecyclable resetWords(long[] words) { return this; } - private final AtomicBoolean recycled = new AtomicBoolean(false); private Handle recyclerHandle = null; private static final Recycler RECYCLER = new Recycler() { @@ -1199,13 +1198,11 @@ private BitSetRecyclable(Handle recyclerHandle) { } public static BitSetRecyclable create() { - BitSetRecyclable instance = RECYCLER.get(); - instance.recycled.set(false); - return instance; + return RECYCLER.get(); } public void recycle() { - if (recyclerHandle != null && recycled.compareAndSet(false, true)) { + if (recyclerHandle != null) { this.clear(); recyclerHandle.recycle(this); } From 7b72bf7880c3b9c000d74ff8bdda1ff892e0c6dc Mon Sep 17 00:00:00 2001 From: lipenghui Date: Fri, 3 Jul 2020 09:06:54 +0800 Subject: [PATCH 3/3] Update pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java --- .../apache/pulsar/common/util/collections/BitSetRecyclable.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java index c586a54d46f77..285416ac606b2 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/BitSetRecyclable.java @@ -27,7 +27,6 @@ import java.nio.LongBuffer; import java.util.Arrays; import java.util.BitSet; -import java.util.concurrent.atomic.AtomicBoolean; /** * This this copy of {@link BitSet}.