From 02117887819b2232e85c4dbfdac5eebb496cc12f Mon Sep 17 00:00:00 2001 From: labuladong Date: Sun, 29 Jan 2023 15:25:45 +0800 Subject: [PATCH] fix stats --- .../pulsar/broker/service/Consumer.java | 5 ++++ .../broker/stats/ConsumerStatsTest.java | 26 +++++++++++++++++++ 2 files changed, 31 insertions(+) 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 47da95b34ac3f..007854a73cc4d 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 @@ -473,6 +473,11 @@ private CompletableFuture individualAckNormal(CommandAck ack, Map 0) { long[] ackSets = new long[msgId.getAckSetsCount()]; for (int j = 0; j < msgId.getAckSetsCount(); j++) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java index fd552766569fe..9ca7c509558c0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java @@ -457,4 +457,30 @@ public void testAvgMessagesPerEntry() throws Exception { int avgMessagesPerEntry = consumerStats.getAvgMessagesPerEntry(); assertEquals(3, avgMessagesPerEntry); } + + + @Test + public void testDuplicateAcknowledgement() + throws PulsarClientException, PulsarAdminException { + final String topicName = "persistent://public/default/duplicated-acknowledgement-test"; + @Cleanup + Producer producer1 = pulsarClient.newProducer() + .topic(topicName) + .create(); + @Cleanup + Consumer consumer1 = pulsarClient.newConsumer() + .topic(topicName) + .subscriptionName("sub-1") + .acknowledgmentGroupTime(0, TimeUnit.SECONDS) + .subscriptionType(SubscriptionType.Shared) + .subscribe(); + producer1.send("1".getBytes()); + Message message = consumer1.receive(); + assertEquals(admin.topics().getStats(topicName).getSubscriptions() + .get("sub-1").getUnackedMessages(), 1); + consumer1.acknowledge(message); + consumer1.acknowledge(message); + assertEquals(admin.topics().getStats(topicName).getSubscriptions() + .get("sub-1").getUnackedMessages(), 0); + } }