From 7a804cd877107bff94fd493f81da12a4a5f2e060 Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Wed, 1 Jun 2022 16:24:17 +0800 Subject: [PATCH 1/3] [fix][txn]Fix transasction ack batch message --- .../pulsar/broker/service/Consumer.java | 2 +- .../pendingack/PendingAckPersistentTest.java | 71 +++++++++++++++++++ 2 files changed, 72 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 8dfb7a3345e2f..c8883982c4114 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -478,7 +478,7 @@ private CompletableFuture individualAckWithTransaction(CommandAck ack) { MessageIdData msgId = ack.getMessageIdAt(i); PositionImpl position; long ackedCount = 0; - long batchSize = getBatchSize(msgId); + long batchSize = msgId.getBatchSize(); Consumer ackOwnerConsumer = getAckOwnerConsumer(msgId.getLedgerId(), msgId.getEntryId()); if (msgId.getAckSetsCount() > 0) { long[] ackSets = new long[msgId.getAckSetsCount()]; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java index 96b2b9c14bfec..241076b93bd41 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java @@ -46,8 +46,10 @@ import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.transaction.Transaction; +import org.apache.pulsar.client.impl.BatchMessageIdImpl; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.client.impl.transaction.TransactionImpl; @@ -567,4 +569,73 @@ public void testPendingAckLowWaterMarkRemoveFirstTxn() throws Exception { assertFalse(individualAckOfTransaction.containsKey(transaction2.getTxnID())); } + + @Test + public void testTransactionConflictExceptionWhenAckBatchMessage() throws Exception { + String topic = TopicName.get(TopicDomain.persistent.toString(), + NamespaceName.get(NAMESPACE1), "test").toString(); + + String subscriptionName = "my-subscription-batch"; + pulsarServiceList.get(0).getBrokerService() + .getManagedLedgerConfig(TopicName.get(topic)).get() + .setDeletionAtBatchIndexLevelEnabled(true); + + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .enableBatching(true) + .batchingMaxMessages(3) + // set batch max publish delay big enough to make sure entry has 3 messages + .batchingMaxPublishDelay(10, TimeUnit.SECONDS) + .topic(topic).create(); + + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .subscriptionName(subscriptionName) + .enableBatchIndexAcknowledgment(true) + .subscriptionType(SubscriptionType.Exclusive) + .isAckReceiptEnabled(true) + .topic(topic) + .subscribe(); + + List messageIds = new ArrayList<>(); + List> futureMessageIds = new ArrayList<>(); + + List messages = new ArrayList<>(); + for (int i = 0; i < 3; i++) { + String message = "my-message-" + i; + messages.add(message); + CompletableFuture messageIdCompletableFuture = producer.sendAsync(message); + futureMessageIds.add(messageIdCompletableFuture); + } + + for (CompletableFuture futureMessageId : futureMessageIds) { + MessageId messageId = futureMessageId.get(); + messageIds.add(messageId); + } + + Transaction transaction = pulsarClient.newTransaction() + .withTransactionTimeout(5, TimeUnit.DAYS) + .build() + .get(); + + Message message1 = consumer.receive(); + Message message2 = consumer.receive(); + + BatchMessageIdImpl messageId = (BatchMessageIdImpl) message2.getMessageId(); + consumer.acknowledgeAsync(messageId, transaction).get(); + + Transaction transaction2 = pulsarClient.newTransaction() + .withTransactionTimeout(5, TimeUnit.DAYS) + .build() + .get(); + transaction.commit().get(); + + try { + consumer.acknowledgeAsync(messageId, transaction2).get(); + fail(); + } catch (ExecutionException e) { + Assert.assertTrue(e.getCause() instanceof PulsarClientException.TransactionConflictException); + } + } + } From 92dbf03ee8f9c60aeeb7ea2d5f4b4856c9a63e2d Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Thu, 2 Jun 2022 19:49:56 +0800 Subject: [PATCH 2/3] optimize --- .../main/java/org/apache/pulsar/broker/service/Consumer.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index c8883982c4114..44f1f8fb8357a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -478,7 +478,7 @@ private CompletableFuture individualAckWithTransaction(CommandAck ack) { MessageIdData msgId = ack.getMessageIdAt(i); PositionImpl position; long ackedCount = 0; - long batchSize = msgId.getBatchSize(); + long batchSize = getBatchSize(msgId); Consumer ackOwnerConsumer = getAckOwnerConsumer(msgId.getLedgerId(), msgId.getEntryId()); if (msgId.getAckSetsCount() > 0) { long[] ackSets = new long[msgId.getAckSetsCount()]; @@ -492,7 +492,7 @@ private CompletableFuture individualAckWithTransaction(CommandAck ack) { ackedCount = batchSize; } - positionsAcked.add(new MutablePair<>(position, (int) batchSize)); + positionsAcked.add(new MutablePair<>(position, msgId.getBatchSize())); addAndGetUnAckedMsgs(ackOwnerConsumer, -(int) ackedCount); From 69b19bb4042dd343d36bc01b46a4ef2ebfe794ac Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Tue, 7 Jun 2022 09:08:56 +0800 Subject: [PATCH 3/3] optimize --- .../java/org/apache/pulsar/broker/service/Consumer.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index e293df9b5922a..161e2e7e7f8ba 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -493,8 +493,11 @@ private CompletableFuture individualAckWithTransaction(CommandAck ack) { position = PositionImpl.get(msgId.getLedgerId(), msgId.getEntryId()); ackedCount = batchSize; } - - positionsAcked.add(new MutablePair<>(position, msgId.getBatchSize())); + if (msgId.hasBatchSize()) { + positionsAcked.add(new MutablePair<>(position, msgId.getBatchSize())); + } else { + positionsAcked.add(new MutablePair<>(position, (int) batchSize)); + } addAndGetUnAckedMsgs(ackOwnerConsumer, -(int) ackedCount);