diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index 232bf5bd2ca89..bda3f973ebcaa 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -1122,6 +1122,16 @@ public long getEstimatedBacklogSize() { } } + /** + * Get estimated backlog size from a specific position. + */ + public long getEstimatedBacklogSize(PositionImpl pos) { + if (pos == null) { + return 0; + } + return estimateBacklogFromPosition(pos); + } + long estimateBacklogFromPosition(PositionImpl pos) { synchronized (this) { LedgerInfo ledgerInfo = ledgers.get(pos.getLedgerId()); diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 34746f6cf6123..7820d61c0d00d 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -1839,6 +1839,13 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private boolean exposePreciseBacklogInPrometheus = false; + @FieldContext( + category = CATEGORY_METRICS, + doc = "Enable expose the backlog size for each subscription when generating stats.\n" + + " Locking is used for fetching the status so default to false." + ) + private boolean exposeSubscriptionBacklogSizeInPrometheus = false; + /**** --- Functions --- ****/ @FieldContext( category = CATEGORY_FUNCTIONS, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 30fe17f5140ee..21ac617bd856a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -1128,14 +1128,15 @@ private void internalGetSubscriptionsForNonPartitionedTopic(AsyncResponse asyncR } } - protected TopicStats internalGetStats(boolean authoritative, boolean getPreciseBacklog) { + protected TopicStats internalGetStats(boolean authoritative, boolean getPreciseBacklog, + boolean subscriptionBacklogSize) { validateAdminAndClientPermission(); if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); } validateTopicOwnership(topicName, authoritative); Topic topic = getTopicReference(topicName); - return topic.getStats(getPreciseBacklog); + return topic.getStats(getPreciseBacklog, subscriptionBacklogSize); } protected PersistentTopicInternalStats internalGetInternalStats(boolean authoritative, boolean metadata) { @@ -1256,7 +1257,7 @@ public void getInfoFailed(ManagedLedgerException exception, Object ctx) { } protected void internalGetPartitionedStats(AsyncResponse asyncResponse, boolean authoritative, - boolean perPartition, boolean getPreciseBacklog) { + boolean perPartition, boolean getPreciseBacklog, boolean subscriptionBacklogSize) { if (topicName.isGlobal()) { try { validateGlobalNamespaceOwnership(namespaceName); @@ -1278,7 +1279,8 @@ protected void internalGetPartitionedStats(AsyncResponse asyncResponse, boolean try { topicStatsFutureList .add(pulsar().getAdminClient().topics().getStatsAsync( - (topicName.getPartition(i).toString()), getPreciseBacklog)); + (topicName.getPartition(i).toString()), getPreciseBacklog, + subscriptionBacklogSize)); } catch (PulsarServerException e) { asyncResponse.resume(new RestException(e)); return; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java index 7c596eac469d0..d9cc845bc6716 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java @@ -100,7 +100,7 @@ public NonPersistentTopicStats getStats(@PathParam("property") String property, validateTopicName(property, cluster, namespace, encodedTopic); validateAdminOperationOnTopic(authoritative); Topic topic = getTopicReference(topicName); - return ((NonPersistentTopic) topic).getStats(false); + return ((NonPersistentTopic) topic).getStats(false, false); } @GET diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java index df1254ce6dee9..78392c286b833 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java @@ -304,7 +304,7 @@ public TopicStats getStats(@PathParam("property") String property, @PathParam("c @PathParam("namespace") String namespace, @PathParam("topic") @Encoded String encodedTopic, @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { validateTopicName(property, cluster, namespace, encodedTopic); - return internalGetStats(authoritative, false); + return internalGetStats(authoritative, false, false); } @GET @@ -352,7 +352,7 @@ public void getPartitionedStats(@Suspended final AsyncResponse asyncResponse, @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { try { validateTopicName(property, cluster, namespace, encodedTopic); - internalGetPartitionedStats(asyncResponse, authoritative, perPartition, false); + internalGetPartitionedStats(asyncResponse, authoritative, perPartition, false, false); } catch (WebApplicationException wae) { asyncResponse.resume(wae); } catch (Exception e) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java index 609b27fd2b6be..cfc6f66a09000 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java @@ -112,12 +112,15 @@ public NonPersistentTopicStats getStats( @PathParam("topic") @Encoded String encodedTopic, @ApiParam(value = "Is authentication required to perform this operation") @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, - @ApiParam(value = "Is return precise backlog or imprecise backlog") - @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog) { + @ApiParam(value = "If return precise backlog or imprecise backlog") + @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog, + @ApiParam(value = "If return backlog size for each subscription, require locking on ledger so be careful " + + "not to use when there's heavy traffic.") + @QueryParam("subscriptionBacklogSize") @DefaultValue("false") boolean subscriptionBacklogSize) { validateTopicName(tenant, namespace, encodedTopic); validateAdminOperationOnTopic(topicName, authoritative); Topic topic = getTopicReference(topicName); - return ((NonPersistentTopic) topic).getStats(getPreciseBacklog); + return ((NonPersistentTopic) topic).getStats(getPreciseBacklog, subscriptionBacklogSize); } @GET diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java index 2d52f0c009c7d..e4b0ddff02b8e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java @@ -930,10 +930,13 @@ public TopicStats getStats( @PathParam("topic") @Encoded String encodedTopic, @ApiParam(value = "Is authentication required to perform this operation") @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, - @ApiParam(value = "Is return precise backlog or imprecise backlog") - @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog) { + @ApiParam(value = "If return precise backlog or imprecise backlog") + @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog, + @ApiParam(value = "If return backlog size for each subscription, require locking on ledger so be careful " + + "not to use when there's heavy traffic.") + @QueryParam("subscriptionBacklogSize") @DefaultValue("false") boolean subscriptionBacklogSize) { validateTopicName(tenant, namespace, encodedTopic); - return internalGetStats(authoritative, getPreciseBacklog); + return internalGetStats(authoritative, getPreciseBacklog, subscriptionBacklogSize); } @GET @@ -1009,11 +1012,15 @@ public void getPartitionedStats( @QueryParam("perPartition") @DefaultValue("true") boolean perPartition, @ApiParam(value = "Is authentication required to perform this operation") @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, - @ApiParam(value = "Is return precise backlog or imprecise backlog") - @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog) { + @ApiParam(value = "If return precise backlog or imprecise backlog") + @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog, + @ApiParam(value = "If return backlog size for each subscription, require locking on ledger so be careful " + + "not to use when there's heavy traffic.") + @QueryParam("subscriptionBacklogSize") @DefaultValue("false") boolean subscriptionBacklogSize) { try { validatePartitionedTopicName(tenant, namespace, encodedTopic); - internalGetPartitionedStats(asyncResponse, authoritative, perPartition, getPreciseBacklog); + internalGetPartitionedStats(asyncResponse, authoritative, perPartition, getPreciseBacklog, + subscriptionBacklogSize); } catch (WebApplicationException wae) { asyncResponse.resume(wae); } catch (Exception e) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index 374e5cb280aff..5988336a951b4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -742,11 +742,11 @@ public long getBytesInCounter() { } public long getMsgOutCounter() { - return getStats(false).msgOutCounter; + return getStats(false, false).msgOutCounter; } public long getBytesOutCounter() { - return getStats(false).bytesOutCounter; + return getStats(false, false).bytesOutCounter; } public boolean isDeleteWhileInactive() { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 6ae3a85eb973f..4e04a6fabf529 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -1731,7 +1731,7 @@ public String generateUniqueProducerName() { public Map getTopicStats() { HashMap stats = new HashMap<>(); - forEachTopic(topic -> stats.put(topic.getName(), topic.getStats(false))); + forEachTopic(topic -> stats.put(topic.getName(), topic.getStats(false, false))); return stats; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index a18e78f36b516..a5db381ea6706 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -194,7 +194,7 @@ void updateRates(NamespaceStats nsStats, NamespaceBundleStats currentBundleStats ConcurrentOpenHashMap getReplicators(); - TopicStats getStats(boolean getPreciseBacklog); + TopicStats getStats(boolean getPreciseBacklog, boolean subscriptionBacklogSize); CompletableFuture getInternalStats(boolean includeLedgerMetadata); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java index 6638db4a5f157..0b1f377f6afec 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java @@ -738,7 +738,7 @@ public void updateRates(NamespaceStats nsStats, NamespaceBundleStats bundleStats } @Override - public NonPersistentTopicStats getStats(boolean getPreciseBacklog) { + public NonPersistentTopicStats getStats(boolean getPreciseBacklog, boolean subscriptionBacklogSize) { NonPersistentTopicStats stats = new NonPersistentTopicStats(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index 19215f4a3a509..2df6b4cd063ad 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -912,7 +912,7 @@ public long estimateBacklogSize() { return cursor.getEstimatedSizeSinceMarkDeletePosition(); } - public SubscriptionStats getStats(Boolean getPreciseBacklog) { + public SubscriptionStats getStats(Boolean getPreciseBacklog, boolean subscriptionBacklogSize) { SubscriptionStats subStats = new SubscriptionStats(); subStats.lastExpireTimestamp = lastExpireTimestamp; subStats.lastConsumedFlowTimestamp = lastConsumedFlowTimestamp; @@ -956,6 +956,10 @@ public SubscriptionStats getStats(Boolean getPreciseBacklog) { } } subStats.msgBacklog = getNumberOfEntriesInBacklog(getPreciseBacklog); + if (subscriptionBacklogSize) { + subStats.backlogSize = ((ManagedLedgerImpl) topic.getManagedLedger()) + .getEstimatedBacklogSize((PositionImpl) cursor.getMarkDeletedPosition()); + } subStats.msgBacklogNoDelayed = subStats.msgBacklog - subStats.msgDelayed; subStats.msgRateExpired = expiryMonitor.getMessageExpiryRate(); subStats.totalMsgExpired = expiryMonitor.getTotalMessageExpired(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 3e9e7de8884f5..cd721e415e820 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -1631,7 +1631,7 @@ public double getLastUpdatedAvgPublishRateInByte() { } @Override - public TopicStats getStats(boolean getPreciseBacklog) { + public TopicStats getStats(boolean getPreciseBacklog, boolean subscriptionBacklogSize) { TopicStats stats = new TopicStats(); @@ -1656,7 +1656,7 @@ public TopicStats getStats(boolean getPreciseBacklog) { stats.waitingPublishers = getWaitingProducersCount(); subscriptions.forEach((name, subscription) -> { - SubscriptionStats subStats = subscription.getStats(getPreciseBacklog); + SubscriptionStats subStats = subscription.getStats(getPreciseBacklog, subscriptionBacklogSize); stats.msgRateOut += subStats.msgRateOut; stats.msgThroughputOut += subStats.msgThroughputOut; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/NamespaceStatsAggregator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/NamespaceStatsAggregator.java index fddd059d0ddf6..2e728a677e520 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/NamespaceStatsAggregator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/NamespaceStatsAggregator.java @@ -63,7 +63,8 @@ public static void generate(PulsarService pulsar, boolean includeTopicMetrics, b bundlesMap.forEach((bundle, topicsMap) -> { topicsMap.forEach((name, topic) -> { getTopicStats(topic, topicStats, includeConsumerMetrics, - pulsar.getConfiguration().isExposePreciseBacklogInPrometheus()); + pulsar.getConfiguration().isExposePreciseBacklogInPrometheus(), + pulsar.getConfiguration().isExposeSubscriptionBacklogSizeInPrometheus()); if (includeTopicMetrics) { topicsCount.add(1); @@ -85,7 +86,7 @@ public static void generate(PulsarService pulsar, boolean includeTopicMetrics, b } private static void getTopicStats(Topic topic, TopicStats stats, boolean includeConsumerMetrics, - boolean getPreciseBacklog) { + boolean getPreciseBacklog, boolean subscriptionBacklogSize) { stats.reset(); if (topic instanceof PersistentTopic) { @@ -110,7 +111,8 @@ private static void getTopicStats(Topic topic, TopicStats stats, boolean include stats.storageReadRate = mlStats.getReadEntriesRate(); } - org.apache.pulsar.common.policies.data.TopicStats tStatus = topic.getStats(getPreciseBacklog); + org.apache.pulsar.common.policies.data.TopicStats tStatus = topic.getStats(getPreciseBacklog, + subscriptionBacklogSize); stats.msgInCounter = tStatus.msgInCounter; stats.bytesInCounter = tStatus.bytesInCounter; stats.msgOutCounter = tStatus.msgOutCounter; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest2.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest2.java index a79a007951fb1..61f1a1b68dd14 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest2.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest2.java @@ -1154,7 +1154,8 @@ public void testPreciseBacklog() throws PulsarClientException, PulsarAdminExcept TopicStats topicStats = admin.topics().getStats(topic); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 10); - topicStats = admin.topics().getStats(topic, true); + topicStats = admin.topics().getStats(topic, true, true); + assertEquals(topicStats.subscriptions.get(subName).backlogSize, 43); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 1); consumer.acknowledge(message); @@ -1162,7 +1163,8 @@ public void testPreciseBacklog() throws PulsarClientException, PulsarAdminExcept Thread.sleep(500); // Consumer acks the message, so the precise backlog is 0 - topicStats = admin.topics().getStats(topic, true); + topicStats = admin.topics().getStats(topic, true, true); + assertEquals(topicStats.subscriptions.get(subName).backlogSize, 0); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 0); topicStats = admin.topics().getStats(topic); @@ -1207,7 +1209,7 @@ public void testBacklogNoDelayed() throws PulsarClientException, PulsarAdminExce // not yet guaranteed to see the stats updated. Thread.sleep(500); - TopicStats topicStats = admin.topics().getStats(topic, true); + TopicStats topicStats = admin.topics().getStats(topic, true, true); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 10); assertEquals(topicStats.subscriptions.get(subName).msgBacklogNoDelayed, 5); @@ -1216,7 +1218,7 @@ public void testBacklogNoDelayed() throws PulsarClientException, PulsarAdminExce } // Wait the ack send. Thread.sleep(500); - topicStats = admin.topics().getStats(topic, true); + topicStats = admin.topics().getStats(topic, true, true); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 5); assertEquals(topicStats.subscriptions.get(subName).msgBacklogNoDelayed, 0); } @@ -1258,8 +1260,9 @@ public void testPreciseBacklogForPartitionedTopic() throws PulsarClientException TopicStats topicStats = admin.topics().getPartitionedStats(topic, false); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 20); - topicStats = admin.topics().getPartitionedStats(topic, false, true); + topicStats = admin.topics().getPartitionedStats(topic, false, true, true); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 1); + assertEquals(topicStats.subscriptions.get(subName).backlogSize, 43); } @Test(timeOut = 30000) @@ -1296,8 +1299,9 @@ public void testBacklogNoDelayedForPartitionedTopic() throws PulsarClientExcepti } } - TopicStats topicStats = admin.topics().getPartitionedStats(topic, false, true); + TopicStats topicStats = admin.topics().getPartitionedStats(topic, false, true, true); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 10); + assertEquals(topicStats.subscriptions.get(subName).backlogSize, 470); assertEquals(topicStats.subscriptions.get(subName).msgBacklogNoDelayed, 5); for (int i = 0; i < 5; i++) { @@ -1305,8 +1309,9 @@ public void testBacklogNoDelayedForPartitionedTopic() throws PulsarClientExcepti } // Wait the ack send. Thread.sleep(500); - topicStats = admin.topics().getPartitionedStats(topic, false, true); + topicStats = admin.topics().getPartitionedStats(topic, false, true, true); assertEquals(topicStats.subscriptions.get(subName).msgBacklog, 5); + assertEquals(topicStats.subscriptions.get(subName).backlogSize, 238); assertEquals(topicStats.subscriptions.get(subName).msgBacklogNoDelayed, 0); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java index 12dcac5bd30ff..e3a49ab0341d5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java @@ -161,7 +161,7 @@ public void testBrokerServicePersistentTopicStats() throws Exception { assertNotNull(topicRef); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); // subscription stats @@ -182,7 +182,7 @@ public void testBrokerServicePersistentTopicStats() throws Exception { Thread.sleep(ASYNC_EVENT_COMPLETION_WAIT); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); // publisher stats @@ -220,7 +220,7 @@ public void testBrokerServicePersistentTopicStats() throws Exception { Thread.sleep(ASYNC_EVENT_COMPLETION_WAIT); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); assertEquals(stats.offloadedStorageSize, 0); @@ -234,13 +234,13 @@ public void testStatsOfStorageSizeWithSubscription() throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getStats(false).storageSize, 0); + assertEquals(topicRef.getStats(false, false).storageSize, 0); for (int i = 0; i < 10; i++) { producer.send(new byte[10]); } - assertTrue(topicRef.getStats(false).storageSize > 0); + assertTrue(topicRef.getStats(false, false).storageSize > 0); } @Test @@ -259,7 +259,7 @@ public void testBrokerServicePersistentRedeliverTopicStats() throws Exception { assertNotNull(topicRef); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); // subscription stats @@ -277,7 +277,7 @@ public void testBrokerServicePersistentRedeliverTopicStats() throws Exception { Thread.sleep(ASYNC_EVENT_COMPLETION_WAIT); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); // publisher stats @@ -312,7 +312,7 @@ public void testBrokerServicePersistentRedeliverTopicStats() throws Exception { Thread.sleep(ASYNC_EVENT_COMPLETION_WAIT); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); assertTrue(subStats.msgRateRedeliver > 0.0); assertEquals(subStats.msgRateRedeliver, subStats.consumers.get(0).msgRateRedeliver); @@ -326,7 +326,7 @@ public void testBrokerServicePersistentRedeliverTopicStats() throws Exception { Thread.sleep(ASYNC_EVENT_COMPLETION_WAIT); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); assertEquals(subStats.msgBacklog, 0); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DispatcherBlockConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DispatcherBlockConsumerTest.java index 62a6edd6c61c2..604edc3c4f7c9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DispatcherBlockConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DispatcherBlockConsumerTest.java @@ -543,7 +543,7 @@ public void testBlockDispatcherStats() throws Exception { assertNotNull(topicRef); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); // subscription stats @@ -561,7 +561,7 @@ public void testBlockDispatcherStats() throws Exception { Thread.sleep(timeWaitToSync); rolloverPerIntervalStats(); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.subscriptions.values().iterator().next(); assertTrue(subStats.msgBacklog > 0); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/NonPersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/NonPersistentTopicTest.java index 95bd4d9a8f814..bf844d5131e56 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/NonPersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/NonPersistentTopicTest.java @@ -437,7 +437,7 @@ public void testTopicStats() throws Exception { assertNotNull(topicRef); rolloverPerIntervalStats(pulsar); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.getSubscriptions().values().iterator().next(); // subscription stats @@ -455,7 +455,7 @@ public void testTopicStats() throws Exception { Thread.sleep(timeWaitToSync); rolloverPerIntervalStats(pulsar); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.getSubscriptions().values().iterator().next(); assertTrue(subStats.msgRateOut > 0); @@ -519,7 +519,7 @@ public void testReplicator() throws Exception { assertNotNull(replicatorR3); rolloverPerIntervalStats(replicationPulasr); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.getSubscriptions().values().iterator().next(); // subscription stats @@ -590,7 +590,7 @@ public void testReplicator() throws Exception { Thread.sleep(timeWaitToSync); rolloverPerIntervalStats(replicationPulasr); - stats = topicRef.getStats(false); + stats = topicRef.getStats(false, false); subStats = stats.getSubscriptions().values().iterator().next(); assertTrue(subStats.msgRateOut > 0); @@ -809,7 +809,7 @@ public void testMsgDropStat() throws Exception { NonPersistentTopic topic = (NonPersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topicName).get(); pulsar.getBrokerService().updateRates(); - NonPersistentTopicStats stats = topic.getStats(false); + NonPersistentTopicStats stats = topic.getStats(false, false); NonPersistentPublisherStats npStats = stats.getPublishers().get(0); NonPersistentSubscriptionStats sub1Stats = stats.getSubscriptions().get("subscriber-1"); NonPersistentSubscriptionStats sub2Stats = stats.getSubscriptions().get("subscriber-2"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java index 1d3af92717dfd..2f6a02b46a284 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java @@ -139,7 +139,7 @@ public void testLargeMessage(boolean ackReceiptEnabled) throws Exception { pulsar.getBrokerService().updateRates(); - PublisherStats producerStats = topic.getStats(false).publishers.get(0); + PublisherStats producerStats = topic.getStats(false, false).publishers.get(0); assertTrue(producerStats.chunkedMessageRate > 0); diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java index 04178f8760766..738543c32db4f 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java @@ -755,6 +755,8 @@ default CompletableFuture deleteAsync(String topic, boolean force) { * topic name * @param getPreciseBacklog * Set to true to get precise backlog, Otherwise get imprecise backlog. + * @param subscriptionBacklogSize + * Whether to get backlog size for each subscription. * @return the topic statistics * * @throws NotAuthorizedException @@ -764,10 +766,15 @@ default CompletableFuture deleteAsync(String topic, boolean force) { * @throws PulsarAdminException * Unexpected error */ - TopicStats getStats(String topic, boolean getPreciseBacklog) throws PulsarAdminException; + TopicStats getStats(String topic, boolean getPreciseBacklog, + boolean subscriptionBacklogSize) throws PulsarAdminException; + + default TopicStats getStats(String topic, boolean getPreciseBacklog) throws PulsarAdminException { + return getStats(topic, getPreciseBacklog, false); + } default TopicStats getStats(String topic) throws PulsarAdminException { - return getStats(topic, false); + return getStats(topic, false, false); } /** @@ -778,14 +785,16 @@ default TopicStats getStats(String topic) throws PulsarAdminException { * topic name * @param getPreciseBacklog * Set to true to get precise backlog, Otherwise get imprecise backlog. - * + * @param subscriptionBacklogSize + * Whether to get backlog size for each subscription. * @return a future that can be used to track when the topic statistics are returned * */ - CompletableFuture getStatsAsync(String topic, boolean getPreciseBacklog); + CompletableFuture getStatsAsync(String topic, boolean getPreciseBacklog, + boolean subscriptionBacklogSize); default CompletableFuture getStatsAsync(String topic) { - return getStatsAsync(topic, false); + return getStatsAsync(topic, false, false); } /** @@ -936,7 +945,11 @@ default CompletableFuture getStatsAsync(String topic) { * @param topic * topic name * @param perPartition - * + * flag to get stats per partition + * @param getPreciseBacklog + * Set to true to get precise backlog, Otherwise get imprecise backlog. + * @param subscriptionBacklogSize + * Whether to get backlog size for each subscription. * @return the partitioned topic statistics * @throws NotAuthorizedException * Don't have admin permission @@ -946,11 +959,12 @@ default CompletableFuture getStatsAsync(String topic) { * Unexpected error * */ - PartitionedTopicStats getPartitionedStats(String topic, boolean perPartition, boolean getPreciseBacklog) + PartitionedTopicStats getPartitionedStats(String topic, boolean perPartition, boolean getPreciseBacklog, + boolean subscriptionBacklogSize) throws PulsarAdminException; default PartitionedTopicStats getPartitionedStats(String topic, boolean perPartition) throws PulsarAdminException { - return getPartitionedStats(topic, perPartition, false); + return getPartitionedStats(topic, perPartition, false, false); } /** @@ -960,13 +974,17 @@ default PartitionedTopicStats getPartitionedStats(String topic, boolean perParti * topic Name * @param perPartition * flag to get stats per partition + * @param getPreciseBacklog + * Set to true to get precise backlog, Otherwise get imprecise backlog. + * @param subscriptionBacklogSize + * Whether to get backlog size for each subscription. * @return a future that can be used to track when the partitioned topic statistics are returned */ CompletableFuture getPartitionedStatsAsync( - String topic, boolean perPartition, boolean getPreciseBacklog); + String topic, boolean perPartition, boolean getPreciseBacklog, boolean subscriptionBacklogSize); default CompletableFuture getPartitionedStatsAsync(String topic, boolean perPartition) { - return getPartitionedStatsAsync(topic, perPartition, false); + return getPartitionedStatsAsync(topic, perPartition, false, false); } /** diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java index 2a27382a09c84..7d62ae168afbe 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java @@ -626,9 +626,11 @@ public void failed(Throwable throwable) { } @Override - public TopicStats getStats(String topic, boolean getPreciseBacklog) throws PulsarAdminException { + public TopicStats getStats(String topic, boolean getPreciseBacklog, + boolean subscriptionBacklogSize) throws PulsarAdminException { try { - return getStatsAsync(topic, getPreciseBacklog).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); + return getStatsAsync(topic, getPreciseBacklog, subscriptionBacklogSize) + .get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { @@ -640,9 +642,12 @@ public TopicStats getStats(String topic, boolean getPreciseBacklog) throws Pulsa } @Override - public CompletableFuture getStatsAsync(String topic, boolean getPreciseBacklog) { + public CompletableFuture getStatsAsync(String topic, boolean getPreciseBacklog, + boolean subscriptionBacklogSize) { TopicName tn = validateTopic(topic); - WebTarget path = topicPath(tn, "stats").queryParam("getPreciseBacklog", getPreciseBacklog); + WebTarget path = topicPath(tn, "stats") + .queryParam("getPreciseBacklog", getPreciseBacklog) + .queryParam("subscriptionBacklogSize", subscriptionBacklogSize); final CompletableFuture future = new CompletableFuture<>(); asyncGetRequest(path, new InvocationCallback() { @@ -742,10 +747,11 @@ public void failed(Throwable throwable) { } @Override - public PartitionedTopicStats getPartitionedStats(String topic, boolean perPartition, boolean getPreciseBacklog) + public PartitionedTopicStats getPartitionedStats(String topic, boolean perPartition, boolean getPreciseBacklog, + boolean subscriptionBacklogSize) throws PulsarAdminException { try { - return getPartitionedStatsAsync(topic, perPartition, getPreciseBacklog) + return getPartitionedStatsAsync(topic, perPartition, getPreciseBacklog, subscriptionBacklogSize) .get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); @@ -759,10 +765,12 @@ public PartitionedTopicStats getPartitionedStats(String topic, boolean perPartit @Override public CompletableFuture getPartitionedStatsAsync(String topic, - boolean perPartition, boolean getPreciseBacklog) { + boolean perPartition, boolean getPreciseBacklog, boolean subscriptionBacklogSize) { TopicName tn = validateTopic(topic); WebTarget path = topicPath(tn, "partitioned-stats"); - path = path.queryParam("perPartition", perPartition).queryParam("getPreciseBacklog", getPreciseBacklog); + path = path.queryParam("perPartition", perPartition) + .queryParam("getPreciseBacklog", getPreciseBacklog) + .queryParam("subscriptionBacklogSize", subscriptionBacklogSize); final CompletableFuture future = new CompletableFuture<>(); asyncGetRequest(path, new InvocationCallback() { diff --git a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java index cf794e0f1695d..cbf25728f4f86 100644 --- a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java +++ b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java @@ -722,7 +722,7 @@ public void topics() throws Exception { verify(mockTopics).deleteSubscription("persistent://myprop/clust/ns1/ds1", "sub1", false); cmdTopics.run(split("stats persistent://myprop/clust/ns1/ds1")); - verify(mockTopics).getStats("persistent://myprop/clust/ns1/ds1", false); + verify(mockTopics).getStats("persistent://myprop/clust/ns1/ds1", false, false); cmdTopics.run(split("stats-internal persistent://myprop/clust/ns1/ds1")); verify(mockTopics).getInternalStats("persistent://myprop/clust/ns1/ds1", false); @@ -731,7 +731,7 @@ public void topics() throws Exception { verify(mockTopics).getInternalInfo("persistent://myprop/clust/ns1/ds1"); cmdTopics.run(split("partitioned-stats persistent://myprop/clust/ns1/ds1 --per-partition")); - verify(mockTopics).getPartitionedStats("persistent://myprop/clust/ns1/ds1", true, false); + verify(mockTopics).getPartitionedStats("persistent://myprop/clust/ns1/ds1", true, false, false); cmdTopics.run(split("partitioned-stats-internal persistent://myprop/clust/ns1/ds1")); verify(mockTopics).getPartitionedInternalStats("persistent://myprop/clust/ns1/ds1"); diff --git a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java index 59f29c5eac47d..490f0551b82a0 100644 --- a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java +++ b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java @@ -534,10 +534,15 @@ private class GetStats extends CliCommand { "--get-precise-backlog" }, description = "Set true to get precise backlog") private boolean getPreciseBacklog = false; + @Parameter(names = { "-sbs", + "--get-subscription-backlog-size" }, description = "Set true to get backlog size for each subscription" + + ", locking required.") + private boolean subscriptionBacklogSize = false; + @Override void run() throws PulsarAdminException { String topic = validateTopicName(params); - print(getTopics().getStats(topic, getPreciseBacklog)); + print(getTopics().getStats(topic, getPreciseBacklog, subscriptionBacklogSize)); } } @@ -584,10 +589,15 @@ private class GetPartitionedStats extends CliCommand { "--get-precise-backlog" }, description = "Set true to get precise backlog") private boolean getPreciseBacklog = false; + @Parameter(names = { "-sbs", + "--get-subscription-backlog-size" }, description = "Set true to get backlog size for each subscription" + + ", locking required.") + private boolean subscriptionBacklogSize = false; + @Override void run() throws Exception { String topic = validateTopicName(params); - print(getTopics().getPartitionedStats(topic, perPartition, getPreciseBacklog)); + print(topics.getPartitionedStats(topic, perPartition, getPreciseBacklog, subscriptionBacklogSize)); } } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/SubscriptionStats.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/SubscriptionStats.java index 38aec473f437c..13862e160b1ab 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/SubscriptionStats.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/SubscriptionStats.java @@ -50,6 +50,9 @@ public class SubscriptionStats { /** Number of messages in the subscription backlog. */ public long msgBacklog; + /** Size of backlog in byte. **/ + public long backlogSize; + /** Number of messages in the subscription backlog that do not contain the delay messages. */ public long msgBacklogNoDelayed; @@ -119,6 +122,7 @@ public void reset() { msgOutCounter = 0; msgRateRedeliver = 0; msgBacklog = 0; + backlogSize = 0; msgBacklogNoDelayed = 0; unackedMessages = 0; msgRateExpired = 0; @@ -141,6 +145,7 @@ public SubscriptionStats add(SubscriptionStats stats) { this.msgOutCounter += stats.msgOutCounter; this.msgRateRedeliver += stats.msgRateRedeliver; this.msgBacklog += stats.msgBacklog; + this.backlogSize += stats.backlogSize; this.msgBacklogNoDelayed += stats.msgBacklogNoDelayed; this.unackedMessages += stats.unackedMessages; this.msgRateExpired += stats.msgRateExpired;