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 a9831e2d35f11..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, (int) batchSize)); + if (msgId.hasBatchSize()) { + positionsAcked.add(new MutablePair<>(position, msgId.getBatchSize())); + } else { + positionsAcked.add(new MutablePair<>(position, (int) batchSize)); + } addAndGetUnAckedMsgs(ackOwnerConsumer, -(int) ackedCount); 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 eb2cf0955c683..594eb40c5c63b 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; @@ -574,4 +576,73 @@ public void testPendingAckLowWaterMarkRemoveFirstTxn() throws Exception { assertFalse(individualAckOfTransaction.containsKey(transaction1.getTxnID())); 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); + } + } + }