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..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 @@ -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 testDoNotRecycleAckSetMultipleTimes() 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-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();