From ce85df545ff0fa9447471bd60abedaac8d5f7e4e Mon Sep 17 00:00:00 2001 From: mayozhang Date: Tue, 21 Feb 2023 19:03:51 +0800 Subject: [PATCH 1/4] fix receive duplicated message due to pendingAcks in PendingAckHandle --- .../mledger/util/PositionAckSetUtil.java | 5 ++ .../service/AbstractBaseDispatcher.java | 16 ++++ .../client/impl/TransactionEndToEndTest.java | 78 +++++++++++++++++++ .../BitSetRecyclableRecyclableTest.java | 17 ++++ 4 files changed, 116 insertions(+) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java index 8173b30c4fea9..0c40d22eb925c 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java @@ -59,6 +59,11 @@ public static long[] andAckSet(long[] firstAckSet, long[] secondAckSet) { return ackSet; } + public static boolean isAckSetEmpty(long[] ackSet) { + BitSetRecyclable bitSet = BitSetRecyclable.create().resetWords(ackSet); + return bitSet.isEmpty(); + } + //This method is compare two position which position is bigger than another one. //When the ledgerId and entryId in this position is same to another one and two position all have ack set, it will //compare the ack set next bit index is bigger than another one. diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index ef2fd80302a98..4f543f845b139 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -19,6 +19,7 @@ package org.apache.pulsar.broker.service; import static org.apache.bookkeeper.mledger.util.PositionAckSetUtil.andAckSet; +import static org.apache.bookkeeper.mledger.util.PositionAckSetUtil.isAckSetEmpty; import io.netty.buffer.ByteBuf; import io.prometheus.client.Gauge; import java.util.ArrayList; @@ -239,6 +240,21 @@ public int filterEntriesForConsumer(@Nullable MessageMetadata[] metadataArray, i // if actSet is null, use pendingAck ackSet ackSet = positionInPendingAck.getAckSet(); } + // if the result of pendingAckSet(in pendingAckHandle) AND the ackSet(in cursor) is empty + // filter this entry + if (isAckSetEmpty(ackSet)) { + entries.set(i, null); + entry.release(); + continue; + } + } else { + // filter non-batch message in pendingAck state + if (positionInPendingAck.getLedgerId() == entry.getLedgerId() + && positionInPendingAck.getEntryId() == entry.getEntryId()) { + entries.set(i, null); + entry.release(); + continue; + } } } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java index 527b8532e0452..83feaa3ac1158 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java @@ -176,6 +176,84 @@ private void testIndividualAckAbortFilterAckSetInPendingAckState() throws Except assertNull(consumer.receive(2, TimeUnit.SECONDS)); } + + @Test(dataProvider="enableBatch") + private void testFilterMsgsInPendingAckStateWhenConsumerDisconnect(boolean enableBatch) throws Exception { + final String topicName = NAMESPACE1 + "/testFilterMsgsInPendingAckStateWhenConsumerDisconnect-" + enableBatch; + final int count = 10; + + @Cleanup + Producer producer = null; + if (enableBatch) { + producer = pulsarClient + .newProducer(Schema.INT32) + .topic(topicName) + .enableBatching(true) + .batchingMaxPublishDelay(1, TimeUnit.HOURS) + .batchingMaxMessages(count).create(); + } else { + producer = pulsarClient + .newProducer(Schema.INT32) + .topic(topicName) + .enableBatching(false).create(); + } + + @Cleanup + Consumer consumer = pulsarClient + .newConsumer(Schema.INT32) + .topic(topicName) + .isAckReceiptEnabled(true) + .subscriptionName("test") + .subscriptionType(SubscriptionType.Shared) + .enableBatchIndexAcknowledgment(true) + .subscribe(); + + for (int i = 0; i < count; i++) { + producer.sendAsync(i); + } + + Transaction txn1 = getTxn(); + + Transaction txn2 = getTxn(); + + + // txn1 ack half of messages and don't end the txn1 + for (int i = 0; i < count / 2; i++) { + consumer.acknowledgeAsync(consumer.receive().getMessageId(), txn1).get(); + } + + // txn2 ack the rest half of messages and commit tnx2 + for (int i = count / 2; i < count; i++) { + consumer.acknowledgeAsync(consumer.receive().getMessageId(), txn2).get(); + } + // commit txn2 + txn2.commit().get(); + + // close and re-create consumer + consumer.close(); + consumer = pulsarClient + .newConsumer(Schema.INT32) + .topic(topicName) + .isAckReceiptEnabled(true) + .subscriptionName("test") + .subscriptionType(SubscriptionType.Shared) + .enableBatchIndexAcknowledgment(true) + .subscribe(); + + Message message = consumer.receive(3, TimeUnit.SECONDS); + Assert.assertNull(message); + + // abort txn1 + txn1.abort().get(); + // after txn1 aborted, consumer will receive messages txn1 contains + int receiveCounter = 0; + while((message = consumer.receive(3, TimeUnit.SECONDS)) != null) { + Assert.assertEquals(message.getValue().intValue(), receiveCounter); + receiveCounter ++; + } + Assert.assertEquals(receiveCounter, count / 2); + } + @Test(dataProvider="enableBatch") private void produceCommitTest(boolean enableBatch) throws Exception { @Cleanup diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/BitSetRecyclableRecyclableTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/BitSetRecyclableRecyclableTest.java index 8374f2db8961a..8061f853d66c1 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/BitSetRecyclableRecyclableTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/BitSetRecyclableRecyclableTest.java @@ -45,4 +45,21 @@ public void testResetWords() { Assert.assertTrue(bitset1.get(128)); Assert.assertFalse(bitset1.get(256)); } + + @Test + public void testBitSetEmpty() { + BitSetRecyclable bitSet = BitSetRecyclable.create(); + bitSet.set(0, 5); + bitSet.clear(1); + bitSet.clear(2); + bitSet.clear(3); + long[] array = bitSet.toLongArray(); + Assert.assertFalse(bitSet.isEmpty()); + Assert.assertFalse(BitSetRecyclable.create().resetWords(array).isEmpty()); + bitSet.clear(0); + bitSet.clear(4); + Assert.assertTrue(bitSet.isEmpty()); + long[] array1 = bitSet.toLongArray(); + Assert.assertTrue(BitSetRecyclable.create().resetWords(array1).isEmpty()); + } } From 64e64f9283c29b85403f39d5f9fe66dd9d05a719 Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Tue, 21 Feb 2023 22:22:55 +0800 Subject: [PATCH 2/4] recycle bitSet --- .../org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java | 1 + 1 file changed, 1 insertion(+) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java index 0c40d22eb925c..7a4db83f34765 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java @@ -61,6 +61,7 @@ public static long[] andAckSet(long[] firstAckSet, long[] secondAckSet) { public static boolean isAckSetEmpty(long[] ackSet) { BitSetRecyclable bitSet = BitSetRecyclable.create().resetWords(ackSet); + bitSet.recycle(); return bitSet.isEmpty(); } From 84f42a39a16f39558e558b991b4fc903026bc082 Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Tue, 21 Feb 2023 23:02:31 +0800 Subject: [PATCH 3/4] recycle bitset --- .../org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java index 7a4db83f34765..1c607582076a8 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/PositionAckSetUtil.java @@ -61,8 +61,9 @@ public static long[] andAckSet(long[] firstAckSet, long[] secondAckSet) { public static boolean isAckSetEmpty(long[] ackSet) { BitSetRecyclable bitSet = BitSetRecyclable.create().resetWords(ackSet); + boolean isEmpty = bitSet.isEmpty(); bitSet.recycle(); - return bitSet.isEmpty(); + return isEmpty; } //This method is compare two position which position is bigger than another one. From 2c58abf97b2c76615c779ce6b5d5e8dcaede7a17 Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Wed, 22 Feb 2023 19:54:00 +0800 Subject: [PATCH 4/4] apply comment --- .../pulsar/broker/service/AbstractBaseDispatcher.java | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index 4f543f845b139..8f6caa7a20801 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -249,12 +249,9 @@ public int filterEntriesForConsumer(@Nullable MessageMetadata[] metadataArray, i } } else { // filter non-batch message in pendingAck state - if (positionInPendingAck.getLedgerId() == entry.getLedgerId() - && positionInPendingAck.getEntryId() == entry.getEntryId()) { - entries.set(i, null); - entry.release(); - continue; - } + entries.set(i, null); + entry.release(); + continue; } } }