From 16059bba723eddc4c5ca47c7e1ec0aaf9e78f0ce Mon Sep 17 00:00:00 2001 From: ming luo Date: Fri, 3 Jan 2020 23:34:47 -0500 Subject: [PATCH 1/4] prevent redelivery of acked batch message at client api --- .../pulsar/client/impl/NegativeAcksTest.java | 74 +++++++++++ .../pulsar/client/impl/BatchAckedTracker.java | 87 +++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 18 +++ .../client/impl/BatchAckedTrackerTest.java | 117 ++++++++++++++++++ 4 files changed, 296 insertions(+) create mode 100644 pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java create mode 100644 pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/NegativeAcksTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/NegativeAcksTest.java index 857fc20ebd18f..7d899f33c1d51 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/NegativeAcksTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/NegativeAcksTest.java @@ -152,4 +152,78 @@ public void testNegativeAcks(boolean batching, boolean usePartitions, Subscripti consumer.close(); producer.close(); } + + @Test + public void testBatchNegativeAcks() + throws Exception { + + boolean batching = true; + SubscriptionType subscriptionType = SubscriptionType.Shared; + int negAcksDelayMillis = 0; + int ackTimeout = 0; + String topic = "testNegativeAcks-" + System.nanoTime(); + + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topic) + .subscriptionName("sub1") + .acknowledgmentGroupTime(0, TimeUnit.SECONDS) + .subscriptionType(subscriptionType) + .negativeAckRedeliveryDelay(negAcksDelayMillis, TimeUnit.MILLISECONDS) + .ackTimeout(ackTimeout, TimeUnit.MILLISECONDS) + .subscribe(); + + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .enableBatching(batching) + .create(); + + Set nackedMessages = new HashSet<>(); + + final int N = 10; + final int A = 10; + for (int i = 0; i < A; i++) { + String value = "positive-" + i; + producer.sendAsync(value); + } + for (int i = 0; i < N; i++) { + String value = "negative-" + i; + producer.sendAsync(value); + } + producer.flush(); + + // negatively ack negative message + int nacked = 0; + for (int i = 0; i < N + A; i++) { + Message msg = consumer.receive(); + String value = msg.getValue(); + if (value.startsWith("negative")) { + consumer.negativeAcknowledge(msg); + nackedMessages.add(value); + nacked++; + } else { + consumer.acknowledge(msg); + } + } + + assertEquals(nacked, N); + + Set receivedMessages = new HashSet<>(); + + // Only the negatively acknowledged messages should be received again + for (int i = 0; i < N; i++) { + Message msg = consumer.receive(); + receivedMessages.add(msg.getValue()); + consumer.acknowledge(msg); + } + + assertEquals(receivedMessages, nackedMessages); + + // There should be no more messages + assertNull(consumer.receive(100, TimeUnit.MILLISECONDS)); + consumer.close(); + producer.close(); + } + } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java new file mode 100644 index 0000000000000..d5b4ca1ed9ac8 --- /dev/null +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java @@ -0,0 +1,87 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.impl; + +import java.util.HashMap; +import java.util.HashSet; +import java.util.Map; +import java.util.Set; + +import com.google.common.annotations.VisibleForTesting; + +import org.apache.pulsar.client.api.MessageId; + +/** + * Tracks any partial acked batches and its acked messages + * This will prevent acked message redelivery to the client only at the client API level. + * This class does not track batch with all acked message. We trust broker won't deliver again. + */ +class BatchAckedTracker { + + // a map of partial acked batch and its messages already acked + @VisibleForTesting + Map> ackedBatches = new HashMap>(); + + public BatchAckedTracker() { + } + + // we start to track + public void negativeAcked(BatchMessageIdImpl messageId) { + + } + + // If this message should be delivered and tracks this message if it is a batch message + public boolean deliver(MessageId messageId) { + if (messageId instanceof BatchMessageIdImpl) { + BatchMessageIdImpl id = (BatchMessageIdImpl) messageId; + String batchId = getBatchId(id); + if (ackedBatches.containsKey(batchId)){ + Set batch = ackedBatches.get(batchId); + if (batch.contains(id)) { + return false; + } + } + } + // deliver non batch message and any other cases + return true; + } + + /** + * + * @param messageId batchMessageIdImpl + * @return boolean isAllMsgAcked for the batch + */ + public boolean ack (BatchMessageIdImpl messageId) { + String batchId = getBatchId(messageId); + Set batches = ackedBatches.getOrDefault(batchId, new HashSet()); + if (messageId.getBatchSize() == (batches.size() + 1)) { + //we ack complete batch now so delete it from the tracker + ackedBatches.remove(batchId); + return true; + } else { + batches.add(messageId); + ackedBatches.put(batchId, batches); + } + return false; + } + + private static String getBatchId(BatchMessageIdImpl id) { + return id.ledgerId + "-" + id.entryId + "-" + id.partitionIndex; + } +} diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index c7211d2f2bd28..d0a0bfea59e0f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -113,6 +113,7 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle private final ReadWriteLock lock = new ReentrantReadWriteLock(); + private final BatchAckedTracker batchAckedTracker; private final UnAckedMessageTracker unAckedMessageTracker; private final AcknowledgmentsGroupingTracker acknowledgmentsGroupingTracker; private final NegativeAcksTracker negativeAcksTracker; @@ -196,6 +197,7 @@ protected ConsumerImpl(PulsarClientImpl client, String topic, ConsumerConfigurat this.readCompacted = conf.isReadCompacted(); this.subscriptionInitialPosition = conf.getSubscriptionInitialPosition(); this.negativeAcksTracker = new NegativeAcksTracker(this, conf); + this.batchAckedTracker = new BatchAckedTracker(); this.resetIncludeHead = conf.isResetIncludeHead(); this.createTopicIfDoesNotExist = createTopicIfDoesNotExist; @@ -435,6 +437,12 @@ boolean markAckForBatchMessage(BatchMessageIdImpl batchMessageId, AckType ackTyp outstandingAcks = batchMessageId.getOutstandingAcksInSameBatch(); } + if (batchAckedTracker.ack(batchMessageId)) { + // the batch all delievered including previous acked + outstandingAcks = 0; + isAllMsgsAcked = true; + } + int batchSize = batchMessageId.getBatchSize(); // all messages in this batch have been acked if (isAllMsgsAcked) { @@ -1042,6 +1050,16 @@ void receiveIndividualMessagesFromBatch(MessageMetadata msgMetadata, int redeliv BatchMessageIdImpl batchMessageIdImpl = new BatchMessageIdImpl(messageId.getLedgerId(), messageId.getEntryId(), getPartitionIndex(), i, acker); + + if (!batchAckedTracker.deliver(batchMessageIdImpl)) { + // individual batch message has been acked earlier + singleMessagePayload.release(); + singleMessageMetadataBuilder.recycle(); + + ++skippedMessages; + continue; + } + final MessageImpl message = new MessageImpl<>(topicName.toString(), batchMessageIdImpl, msgMetadata, singleMessageMetadataBuilder.build(), singleMessagePayload, createEncryptionContext(msgMetadata), cnx, schema, redeliveryCount); diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java new file mode 100644 index 0000000000000..e1e9bd523be18 --- /dev/null +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java @@ -0,0 +1,117 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.impl; + +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertTrue; + +import org.apache.pulsar.client.api.MessageId; +import org.testng.annotations.Test; + +public class BatchAckedTrackerTest { + + @Test + public void testDeliveredAndAcked() throws Exception { + BatchAckedTracker tracker = new BatchAckedTracker(); + + MessageId nonBatchMessageId = new MessageIdImpl(1, 1, 1); + + assertTrue(tracker.deliver(nonBatchMessageId)); + assertEquals(tracker.ackedBatches.size(), 0); + + int batchSize = 3; + BatchMessageAcker acker = BatchMessageAcker.newAcker(batchSize); + BatchMessageIdImpl batchMessageId = new BatchMessageIdImpl(1, 1, -1, 2, acker); + assertEquals(batchMessageId.getBatchIndex(), 2); + assertEquals(batchMessageId.getBatchSize(), batchSize); + + // ensure message can be delivered to clients + assertTrue(tracker.deliver(batchMessageId)); + assertEquals(tracker.ackedBatches.size(), 0); + + // ensure the first message ack will be tracked + assertFalse(tracker.ack(batchMessageId)); + assertEquals(tracker.ackedBatches.size(), 1); + + // redeliver an acked message will return false + assertFalse(tracker.deliver(batchMessageId)); + + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 1, acker))); + assertEquals(tracker.ackedBatches.size(), 1); + + // all message are acked in a batch of three messages + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 0, acker))); + assertEquals(tracker.ackedBatches.size(), 0); + + } + + @Test + public void testMultipleBatchDeliveredAndAcked() throws Exception { + BatchAckedTracker tracker = new BatchAckedTracker(); + + int batch1Size = 3; + int batch2Size = 5; + int batch3Size = 8; + + BatchMessageAcker acker1 = BatchMessageAcker.newAcker(batch1Size); + BatchMessageAcker acker2 = BatchMessageAcker.newAcker(batch2Size); + BatchMessageAcker acker3 = BatchMessageAcker.newAcker(batch3Size); + + BatchMessageIdImpl batch1MessageId1 = new BatchMessageIdImpl(1, 1, -1, 0, acker1); + BatchMessageIdImpl batch1MessageId2 = new BatchMessageIdImpl(1, 1, -1, 1, acker1); + + // partial first batch acked + assertTrue(tracker.deliver(batch1MessageId1)); + assertFalse(tracker.ack(batch1MessageId1)); + assertTrue(tracker.deliver(batch1MessageId2)); + assertFalse(tracker.ack(batch1MessageId2)); + assertEquals(tracker.ackedBatches.size(), 1); + + // ensure the first message ack will be tracked + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 0, acker2))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 1, acker2))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 2, acker2))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 3, acker2))); + assertEquals(tracker.ackedBatches.size(), 2); + assertTrue(tracker.deliver(new BatchMessageIdImpl(1, 2, -1, 4, acker2))); + + // redeliver an acked message will return false + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 0, acker3))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 1, acker3))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 2, acker3))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 3, acker3))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 4, acker3))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 5, acker3))); + assertEquals(tracker.ackedBatches.size(), 3); + assertTrue(tracker.deliver(new BatchMessageIdImpl(1, 3, -1, 7, acker3))); + + // all message are acked in a batch of three messages + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 2, acker1))); + assertEquals(tracker.ackedBatches.size(), 2); + + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 4, acker2))); + assertEquals(tracker.ackedBatches.size(), 1); + + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 6, acker3))); + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 7, acker3))); + assertEquals(tracker.ackedBatches.size(), 0); + } + +} From a1c824bb1e65e33bd796b4c4e78293fe6614e887 Mon Sep 17 00:00:00 2001 From: ming luo Date: Mon, 6 Jan 2020 00:28:40 -0500 Subject: [PATCH 2/4] add cumulative ack type --- .../pulsar/client/impl/NegativeAcksTest.java | 66 ++++++++++++++++++ .../pulsar/client/impl/BatchAckedTracker.java | 39 +++++++---- .../pulsar/client/impl/ConsumerImpl.java | 2 +- .../client/impl/BatchAckedTrackerTest.java | 67 +++++++++++++------ 4 files changed, 138 insertions(+), 36 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/NegativeAcksTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/NegativeAcksTest.java index 7d899f33c1d51..c6ef528bdb445 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/NegativeAcksTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/NegativeAcksTest.java @@ -226,4 +226,70 @@ public void testBatchNegativeAcks() producer.close(); } + @Test + public void testBatchNegativeCumulativeAcks() + throws Exception { + boolean batching = true; + SubscriptionType subscriptionType = SubscriptionType.Exclusive; + int negAcksDelayMillis = 0; + int ackTimeout = 1000; + String topic = "testNegativeAcks-" + System.nanoTime(); + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topic) + .subscriptionName("sub1") + .acknowledgmentGroupTime(0, TimeUnit.SECONDS) + .subscriptionType(subscriptionType) + .negativeAckRedeliveryDelay(negAcksDelayMillis, TimeUnit.MILLISECONDS) + .ackTimeout(ackTimeout, TimeUnit.MILLISECONDS) + .subscribe(); + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .enableBatching(batching) + .create(); + Set nackedMessages = new HashSet<>(); + Message msg = null; + final int N = 10; + final int A = 10; + for (int i = 0; i < A; i++) { + String value = "positive-" + i; + producer.sendAsync(value); + } + for (int i = 0; i < N; i++) { + String value = "negative-" + i; + producer.sendAsync(value); + } + producer.flush(); + for (int i = 0; i < A; i++) { + msg = consumer.receive(); + } + // Cumulative ack of messages + consumer.acknowledgeCumulative(msg); + // negatively ack negative message + int nacked = 0; + for (int i = 0; i < N; i++) { + msg = consumer.receive(); + String value = msg.getValue(); + consumer.negativeAcknowledge(msg); + nackedMessages.add(value); + nacked++; + } + assertEquals(nacked, N); + Set receivedMessages = new HashSet<>(); + // Only the negatively acknowledged messages should be received again + for (int i = 0; i < N; i++) { + msg = consumer.receive(); + receivedMessages.add(msg.getValue()); + } + // Cumulative ack of messages + consumer.acknowledgeCumulative(msg); + assertEquals(receivedMessages, nackedMessages); + // Wait for unacked timer to trigger if there are unacked messages + Thread.sleep(1500); + // There should be no more messages since they have all been acked + assertNull(consumer.receive(100, TimeUnit.MILLISECONDS)); + consumer.close(); + producer.close(); + } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java index d5b4ca1ed9ac8..ab5167dcbf61f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java @@ -18,14 +18,17 @@ */ package org.apache.pulsar.client.impl; +import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.Map; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import com.google.common.annotations.VisibleForTesting; import org.apache.pulsar.client.api.MessageId; +import org.apache.pulsar.common.api.proto.PulsarApi.CommandAck.AckType; /** * Tracks any partial acked batches and its acked messages @@ -35,25 +38,21 @@ class BatchAckedTracker { // a map of partial acked batch and its messages already acked + // Key is the string Id for batch, the value is a synchronized set of message Id in string @VisibleForTesting - Map> ackedBatches = new HashMap>(); + Map> ackedBatches = new ConcurrentHashMap>(); public BatchAckedTracker() { } - // we start to track - public void negativeAcked(BatchMessageIdImpl messageId) { - - } - // If this message should be delivered and tracks this message if it is a batch message public boolean deliver(MessageId messageId) { if (messageId instanceof BatchMessageIdImpl) { BatchMessageIdImpl id = (BatchMessageIdImpl) messageId; String batchId = getBatchId(id); if (ackedBatches.containsKey(batchId)){ - Set batch = ackedBatches.get(batchId); - if (batch.contains(id)) { + Set batch = ackedBatches.get(batchId); + if (batch.contains(getBatchMessageId(batchId, id.getBatchIndex()))) { return false; } } @@ -67,21 +66,33 @@ public boolean deliver(MessageId messageId) { * @param messageId batchMessageIdImpl * @return boolean isAllMsgAcked for the batch */ - public boolean ack (BatchMessageIdImpl messageId) { + public boolean ack (BatchMessageIdImpl messageId, AckType ackType) { String batchId = getBatchId(messageId); - Set batches = ackedBatches.getOrDefault(batchId, new HashSet()); - if (messageId.getBatchSize() == (batches.size() + 1)) { + Set batch = ackedBatches.getOrDefault(batchId, + Collections.synchronizedSet(new HashSet())); + if (ackType == AckType.Individual) { + batch.add(getBatchMessageId(batchId, messageId.getBatchIndex())); + } else { + for (int i=0; i<=messageId.getBatchIndex(); i++) { + batch.add(getBatchMessageId(batchId, i)); + } + } + + if (messageId.getBatchSize() == batch.size()) { //we ack complete batch now so delete it from the tracker ackedBatches.remove(batchId); return true; } else { - batches.add(messageId); - ackedBatches.put(batchId, batches); + ackedBatches.put(batchId, batch); + return false; } - return false; } private static String getBatchId(BatchMessageIdImpl id) { return id.ledgerId + "-" + id.entryId + "-" + id.partitionIndex; } + + private static String getBatchMessageId(String batchId, int batchIndex) { + return batchId + "-" + batchIndex; + } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index d0a0bfea59e0f..75c51a7d8befb 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -437,7 +437,7 @@ boolean markAckForBatchMessage(BatchMessageIdImpl batchMessageId, AckType ackTyp outstandingAcks = batchMessageId.getOutstandingAcksInSameBatch(); } - if (batchAckedTracker.ack(batchMessageId)) { + if (batchAckedTracker.ack(batchMessageId, ackType)) { // the batch all delievered including previous acked outstandingAcks = 0; isAllMsgsAcked = true; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java index e1e9bd523be18..3b4406e796fcb 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java @@ -23,6 +23,7 @@ import static org.testng.Assert.assertTrue; import org.apache.pulsar.client.api.MessageId; +import org.apache.pulsar.common.api.proto.PulsarApi.CommandAck.AckType; import org.testng.annotations.Test; public class BatchAckedTrackerTest { @@ -38,8 +39,8 @@ public void testDeliveredAndAcked() throws Exception { int batchSize = 3; BatchMessageAcker acker = BatchMessageAcker.newAcker(batchSize); - BatchMessageIdImpl batchMessageId = new BatchMessageIdImpl(1, 1, -1, 2, acker); - assertEquals(batchMessageId.getBatchIndex(), 2); + BatchMessageIdImpl batchMessageId = new BatchMessageIdImpl(1, 1, -1, 0, acker); + assertEquals(batchMessageId.getBatchIndex(), 0); assertEquals(batchMessageId.getBatchSize(), batchSize); // ensure message can be delivered to clients @@ -47,17 +48,17 @@ public void testDeliveredAndAcked() throws Exception { assertEquals(tracker.ackedBatches.size(), 0); // ensure the first message ack will be tracked - assertFalse(tracker.ack(batchMessageId)); + assertFalse(tracker.ack(batchMessageId, AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 1); // redeliver an acked message will return false assertFalse(tracker.deliver(batchMessageId)); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 1, acker))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 1, acker), AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 1); // all message are acked in a batch of three messages - assertTrue(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 0, acker))); + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 2, acker), AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 0); } @@ -79,39 +80,63 @@ public void testMultipleBatchDeliveredAndAcked() throws Exception { // partial first batch acked assertTrue(tracker.deliver(batch1MessageId1)); - assertFalse(tracker.ack(batch1MessageId1)); + assertFalse(tracker.ack(batch1MessageId1, AckType.Individual)); assertTrue(tracker.deliver(batch1MessageId2)); - assertFalse(tracker.ack(batch1MessageId2)); + assertFalse(tracker.ack(batch1MessageId2, AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 1); // ensure the first message ack will be tracked - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 0, acker2))); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 1, acker2))); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 2, acker2))); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 3, acker2))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 0, acker2), AckType.Individual)); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 1, acker2), AckType.Individual)); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 2, acker2), AckType.Individual)); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 3, acker2), AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 2); assertTrue(tracker.deliver(new BatchMessageIdImpl(1, 2, -1, 4, acker2))); // redeliver an acked message will return false - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 0, acker3))); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 1, acker3))); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 2, acker3))); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 3, acker3))); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 4, acker3))); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 5, acker3))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 0, acker3), AckType.Individual)); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 1, acker3), AckType.Individual)); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 2, acker3), AckType.Individual)); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 3, acker3), AckType.Individual)); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 4, acker3), AckType.Individual)); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 5, acker3), AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 3); assertTrue(tracker.deliver(new BatchMessageIdImpl(1, 3, -1, 7, acker3))); // all message are acked in a batch of three messages - assertTrue(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 2, acker1))); + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 2, acker1), AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 2); - assertTrue(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 4, acker2))); + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 2, -1, 4, acker2), AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 1); - assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 6, acker3))); - assertTrue(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 7, acker3))); + assertFalse(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 6, acker3), AckType.Individual)); + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 3, -1, 7, acker3), AckType.Individual)); assertEquals(tracker.ackedBatches.size(), 0); } + @Test + public void testCumulativeAckedBatch() throws Exception { + BatchAckedTracker tracker = new BatchAckedTracker(); + + int batchSize = 10; + BatchMessageAcker acker = BatchMessageAcker.newAcker(batchSize); + BatchMessageIdImpl batchMessageId = new BatchMessageIdImpl(1, 1, -1, 8, acker); + + // ensure message can be delivered to clients + assertTrue(tracker.deliver(batchMessageId)); + assertEquals(tracker.ackedBatches.size(), 0); + + // all previous messages + assertFalse(tracker.ack(batchMessageId, AckType.Cumulative)); + assertEquals(tracker.ackedBatches.size(), 1); + String batchId = batchMessageId.ledgerId + "-" + batchMessageId.entryId + "-" + batchMessageId.partitionIndex; + assertEquals(tracker.ackedBatches.get(batchId).size(), batchMessageId.getBatchIndex() + 1); + + // ack the batch with the last message + assertTrue(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 9, acker), AckType.Cumulative)); + assertEquals(tracker.ackedBatches.size(), 0); + + } + } From b2410a3b0b318af79245d390c99190cf6a6312ff Mon Sep 17 00:00:00 2001 From: ming luo Date: Mon, 6 Jan 2020 14:27:17 -0500 Subject: [PATCH 3/4] use batch index to track message save memory --- .../pulsar/client/impl/BatchAckedTracker.java | 23 ++++++++----------- 1 file changed, 9 insertions(+), 14 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java index ab5167dcbf61f..1d89c4063c178 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java @@ -19,7 +19,6 @@ package org.apache.pulsar.client.impl; import java.util.Collections; -import java.util.HashMap; import java.util.HashSet; import java.util.Map; import java.util.Set; @@ -38,9 +37,9 @@ class BatchAckedTracker { // a map of partial acked batch and its messages already acked - // Key is the string Id for batch, the value is a synchronized set of message Id in string + // Key is the string Id for batch, the value is a synchronized set of batch index in integer @VisibleForTesting - Map> ackedBatches = new ConcurrentHashMap>(); + Map> ackedBatches = new ConcurrentHashMap>(); public BatchAckedTracker() { } @@ -51,8 +50,8 @@ public boolean deliver(MessageId messageId) { BatchMessageIdImpl id = (BatchMessageIdImpl) messageId; String batchId = getBatchId(id); if (ackedBatches.containsKey(batchId)){ - Set batch = ackedBatches.get(batchId); - if (batch.contains(getBatchMessageId(batchId, id.getBatchIndex()))) { + Set batch = ackedBatches.get(batchId); + if (batch.contains(id.getBatchIndex())) { return false; } } @@ -68,17 +67,17 @@ public boolean deliver(MessageId messageId) { */ public boolean ack (BatchMessageIdImpl messageId, AckType ackType) { String batchId = getBatchId(messageId); - Set batch = ackedBatches.getOrDefault(batchId, - Collections.synchronizedSet(new HashSet())); + Set batch = ackedBatches.getOrDefault(batchId, + Collections.synchronizedSet(new HashSet())); if (ackType == AckType.Individual) { - batch.add(getBatchMessageId(batchId, messageId.getBatchIndex())); + batch.add(messageId.getBatchIndex()); } else { for (int i=0; i<=messageId.getBatchIndex(); i++) { - batch.add(getBatchMessageId(batchId, i)); + batch.add(i); } } - if (messageId.getBatchSize() == batch.size()) { + if (messageId.getBatchSize() <= batch.size()) { //we ack complete batch now so delete it from the tracker ackedBatches.remove(batchId); return true; @@ -91,8 +90,4 @@ public boolean ack (BatchMessageIdImpl messageId, AckType ackType) { private static String getBatchId(BatchMessageIdImpl id) { return id.ledgerId + "-" + id.entryId + "-" + id.partitionIndex; } - - private static String getBatchMessageId(String batchId, int batchIndex) { - return batchId + "-" + batchIndex; - } } From 59e2a1978a549cb857c6b7f5106ea07d9d5b5d64 Mon Sep 17 00:00:00 2001 From: ming luo Date: Tue, 7 Jan 2020 08:55:41 -0500 Subject: [PATCH 4/4] use BitSet to manage batch acked tracker --- .../pulsar/client/impl/BatchAckedTracker.java | 30 ++++++++++--------- .../client/impl/BatchAckedTrackerTest.java | 3 +- 2 files changed, 18 insertions(+), 15 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java index 1d89c4063c178..c2b0c31252d12 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchAckedTracker.java @@ -18,7 +18,9 @@ */ package org.apache.pulsar.client.impl; +import java.util.BitSet; import java.util.Collections; +import java.util.HashMap; import java.util.HashSet; import java.util.Map; import java.util.Set; @@ -37,9 +39,9 @@ class BatchAckedTracker { // a map of partial acked batch and its messages already acked - // Key is the string Id for batch, the value is a synchronized set of batch index in integer + // Key is the string Id for batch, the value is bit set for message index whether ack-ed or not @VisibleForTesting - Map> ackedBatches = new ConcurrentHashMap>(); + Map ackedBatches = Collections.synchronizedMap(new HashMap()); public BatchAckedTracker() { } @@ -49,17 +51,19 @@ public boolean deliver(MessageId messageId) { if (messageId instanceof BatchMessageIdImpl) { BatchMessageIdImpl id = (BatchMessageIdImpl) messageId; String batchId = getBatchId(id); - if (ackedBatches.containsKey(batchId)){ - Set batch = ackedBatches.get(batchId); - if (batch.contains(id.getBatchIndex())) { - return false; - } + if (ackedBatches.containsKey(batchId)) { + return ackedBatches.get(batchId).get(id.getBatchIndex()); } } // deliver non batch message and any other cases return true; } + private BitSet initBatchSet(int size) { + BitSet set = new BitSet(size); + set.set(0, size); + return set; + } /** * * @param messageId batchMessageIdImpl @@ -67,17 +71,15 @@ public boolean deliver(MessageId messageId) { */ public boolean ack (BatchMessageIdImpl messageId, AckType ackType) { String batchId = getBatchId(messageId); - Set batch = ackedBatches.getOrDefault(batchId, - Collections.synchronizedSet(new HashSet())); + int batchSize = messageId.getBatchSize(); + BitSet batch = ackedBatches.getOrDefault(batchId, initBatchSet(batchSize)); if (ackType == AckType.Individual) { - batch.add(messageId.getBatchIndex()); + batch.clear(messageId.getBatchIndex()); } else { - for (int i=0; i<=messageId.getBatchIndex(); i++) { - batch.add(i); - } + batch.clear(0, messageId.getBatchIndex() + 1); } - if (messageId.getBatchSize() <= batch.size()) { + if (batch.isEmpty()) { //we ack complete batch now so delete it from the tracker ackedBatches.remove(batchId); return true; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java index 3b4406e796fcb..a0a0af69923dc 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchAckedTrackerTest.java @@ -131,7 +131,8 @@ public void testCumulativeAckedBatch() throws Exception { assertFalse(tracker.ack(batchMessageId, AckType.Cumulative)); assertEquals(tracker.ackedBatches.size(), 1); String batchId = batchMessageId.ledgerId + "-" + batchMessageId.entryId + "-" + batchMessageId.partitionIndex; - assertEquals(tracker.ackedBatches.get(batchId).size(), batchMessageId.getBatchIndex() + 1); + // assertEquals(tracker.ackedBatches.get(batchId).size(), batchMessageId.getBatchIndex() + 1); + assertEquals(tracker.ackedBatches.get(batchId).nextSetBit(8), 9); // ack the batch with the last message assertTrue(tracker.ack(new BatchMessageIdImpl(1, 1, -1, 9, acker), AckType.Cumulative));