From 7f760e708d5d22de6e14f446a1e65ef041b59023 Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Sun, 24 Jan 2021 14:38:39 -0800 Subject: [PATCH 1/6] Add subscription backlog size info for topicstats. --- .../mledger/impl/ManagedLedgerImpl.java | 10 ++++++++++ .../pulsar/broker/ServiceConfiguration.java | 7 +++++++ .../admin/impl/PersistentTopicsBase.java | 8 ++++---- .../broker/admin/v1/NonPersistentTopics.java | 2 +- .../broker/admin/v1/PersistentTopics.java | 4 ++-- .../broker/admin/v2/NonPersistentTopics.java | 9 ++++++--- .../broker/admin/v2/PersistentTopics.java | 18 ++++++++++++------ .../pulsar/broker/service/AbstractTopic.java | 4 ++-- .../pulsar/broker/service/BrokerService.java | 2 +- .../apache/pulsar/broker/service/Topic.java | 2 +- .../nonpersistent/NonPersistentTopic.java | 2 +- .../persistent/PersistentSubscription.java | 5 ++++- .../service/persistent/PersistentTopic.java | 4 ++-- .../prometheus/NamespaceStatsAggregator.java | 8 +++++--- .../pulsar/broker/admin/AdminApiTest2.java | 10 ++++++---- .../broker/service/BrokerServiceTest.java | 18 +++++++++--------- .../api/DispatcherBlockConsumerTest.java | 4 ++-- .../client/api/NonPersistentTopicTest.java | 10 +++++----- .../client/impl/MessageChunkingTest.java | 2 +- .../org/apache/pulsar/client/admin/Topics.java | 15 ++++++++++----- .../client/admin/internal/TopicsImpl.java | 13 +++++++++---- .../policies/data/SubscriptionStats.java | 5 +++++ 22 files changed, 105 insertions(+), 57 deletions(-) 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 9264e152040bf..251bf132b578a 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 @@ -1117,6 +1117,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 ea457abf9bf92..b70ea8dcea795 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 @@ -1826,6 +1826,13 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private boolean exposePreciseBacklogInPrometheus = false; + @FieldContext( + category = CATEGORY_METRICS, + doc = "Enable expose the backlog size fro each subscription when generating stats.\n" + + " Locking is used for fetching the status so default to false." + ) + private boolean exposeSubscriptionBacklokSizeInPrometheus = 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 edc221e492efe..78b6228cd973e 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 @@ -1109,14 +1109,14 @@ 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) { @@ -1237,7 +1237,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); @@ -1259,7 +1259,7 @@ 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..8738e0dac3141 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 10acb5bce62b2..f7e777b9b12e2 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 @@ -925,10 +925,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 backlob 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 @@ -1004,11 +1007,14 @@ 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 053ee69757e12..bb501ed629350 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 @@ -1730,7 +1730,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 f04b7d655a954..c0c322ec3e7bd 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 @@ -191,7 +191,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 a95692f88a16d..45fd862310292 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 @@ -903,7 +903,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; @@ -947,6 +947,9 @@ 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 cba18613e6ff4..669df625a6ff4 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 @@ -1630,7 +1630,7 @@ public double getLastUpdatedAvgPublishRateInByte() { } @Override - public TopicStats getStats(boolean getPreciseBacklog) { + public TopicStats getStats(boolean getPreciseBacklog, boolean subscriptionBacklogSize) { TopicStats stats = new TopicStats(); @@ -1655,7 +1655,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..562d99c0a9a47 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().isExposeSubscriptionBacklokSizeInPrometheus()); 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 1141e6a977257..3560a468a2fcd 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); } 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 97ddcca0de82c..19fbd98e54487 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 @@ -179,7 +179,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 @@ -216,7 +216,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(subStats.msgBacklog, 0); @@ -229,13 +229,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 @@ -254,7 +254,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 @@ -272,7 +272,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 @@ -307,7 +307,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); @@ -321,7 +321,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 80bf968013b27..f4783f561eb74 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/src/main/java/org/apache/pulsar/client/admin/Topics.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java index 9a4cce48ce460..28b0d497a5679 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java +++ b/pulsar-client-admin/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,11 @@ 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) throws PulsarAdminException { - return getStats(topic, false); + return getStats(topic, false, false); } /** @@ -778,14 +781,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); } /** 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 b5fb762970d5d..628da7616c9dd 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() { 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; From c1f4c353f98f724c64fb37eff2cbb7def44b2a9c Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Sun, 24 Jan 2021 16:32:09 -0800 Subject: [PATCH 2/6] Fixing code style. --- .../admin/impl/PersistentTopicsBase.java | 6 ++++-- .../broker/admin/v2/NonPersistentTopics.java | 4 ++-- .../broker/admin/v2/PersistentTopics.java | 11 ++++++----- .../persistent/PersistentSubscription.java | 3 ++- .../apache/pulsar/client/admin/Topics.java | 19 ++++++++++++++----- .../client/admin/internal/TopicsImpl.java | 11 +++++++---- .../pulsar/admin/cli/PulsarAdminToolTest.java | 2 +- .../apache/pulsar/admin/cli/CmdTopics.java | 14 ++++++++++++-- 8 files changed, 48 insertions(+), 22 deletions(-) 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 78b6228cd973e..d596a059edf6b 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 @@ -1109,7 +1109,8 @@ private void internalGetSubscriptionsForNonPartitionedTopic(AsyncResponse asyncR } } - protected TopicStats internalGetStats(boolean authoritative, boolean getPreciseBacklog, boolean subscriptionBacklogSize) { + protected TopicStats internalGetStats(boolean authoritative, boolean getPreciseBacklog, + boolean subscriptionBacklogSize) { validateAdminAndClientPermission(); if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); @@ -1259,7 +1260,8 @@ protected void internalGetPartitionedStats(AsyncResponse asyncResponse, boolean try { topicStatsFutureList .add(pulsar().getAdminClient().topics().getStatsAsync( - (topicName.getPartition(i).toString()), getPreciseBacklog, subscriptionBacklogSize)); + (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/v2/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java index 8738e0dac3141..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 @@ -114,8 +114,8 @@ public NonPersistentTopicStats getStats( @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, @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.") + @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); 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 f7e777b9b12e2..45c9a0f076459 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 @@ -927,8 +927,8 @@ public TopicStats getStats( @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, @ApiParam(value = "If return precise backlog or imprecise backlog") @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog, - @ApiParam(value = "If return backlob size for each subscription, require locking on ledger so be careful " + - "not to use when there's heavy traffic.") + @ApiParam(value = "If return backlob 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, subscriptionBacklogSize); @@ -1009,12 +1009,13 @@ public void getPartitionedStats( @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, @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.") + @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, subscriptionBacklogSize); + 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/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index 45fd862310292..ec0a1d3baf651 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 @@ -948,7 +948,8 @@ public SubscriptionStats getStats(Boolean getPreciseBacklog, boolean subscriptio } subStats.msgBacklog = getNumberOfEntriesInBacklog(getPreciseBacklog); if (subscriptionBacklogSize) { - subStats.backlogSize = ((ManagedLedgerImpl) topic.getManagedLedger()).getEstimatedBacklogSize((PositionImpl) cursor.getMarkDeletedPosition()); + subStats.backlogSize = ((ManagedLedgerImpl) topic.getManagedLedger()) + .getEstimatedBacklogSize((PositionImpl) cursor.getMarkDeletedPosition()); } subStats.msgBacklogNoDelayed = subStats.msgBacklog - subStats.msgDelayed; subStats.msgRateExpired = expiryMonitor.getMessageExpiryRate(); diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java index 28b0d497a5679..f70d2bee10a7f 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java @@ -941,7 +941,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 @@ -951,11 +955,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); } /** @@ -965,13 +970,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 628da7616c9dd..086c4911f7538 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 @@ -747,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(); @@ -764,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 668f82566c883..8be4deb912205 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 @@ -719,7 +719,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); 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 f565fdde60a8e..2c64bf8087edf 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 @@ -484,10 +484,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(topics.getStats(topic, getPreciseBacklog)); + print(topics.getStats(topic, getPreciseBacklog, subscriptionBacklogSize)); } } @@ -534,10 +539,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(topics.getPartitionedStats(topic, perPartition, getPreciseBacklog)); + print(topics.getPartitionedStats(topic, perPartition, getPreciseBacklog, subscriptionBacklogSize)); } } From 5dde05a46059ef18757834d0a8da87bedd370edc Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Sun, 24 Jan 2021 22:04:18 -0800 Subject: [PATCH 3/6] Fix unit tests. --- .../org/apache/pulsar/broker/admin/AdminApiTest2.java | 9 ++++++--- .../org/apache/pulsar/admin/cli/PulsarAdminToolTest.java | 2 +- 2 files changed, 7 insertions(+), 4 deletions(-) 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 3560a468a2fcd..1a2a190d86e19 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 @@ -1260,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) @@ -1298,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++) { @@ -1307,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-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 8be4deb912205..8c6e17dd926b8 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 @@ -728,7 +728,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"); From 818c6117d0311f6e08dac4561cfcdfe0115d6fc0 Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Mon, 25 Jan 2021 21:14:28 -0800 Subject: [PATCH 4/6] Fix nameing. --- .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 2 +- .../broker/stats/prometheus/NamespaceStatsAggregator.java | 2 +- .../src/main/java/org/apache/pulsar/client/admin/Topics.java | 4 ++++ 3 files changed, 6 insertions(+), 2 deletions(-) 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 b70ea8dcea795..7fe97cc611db1 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 @@ -1831,7 +1831,7 @@ public class ServiceConfiguration implements PulsarConfiguration { doc = "Enable expose the backlog size fro each subscription when generating stats.\n" + " Locking is used for fetching the status so default to false." ) - private boolean exposeSubscriptionBacklokSizeInPrometheus = false; + private boolean exposeSubscriptionBacklogSizeInPrometheus = false; /**** --- Functions --- ****/ @FieldContext( 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 562d99c0a9a47..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 @@ -64,7 +64,7 @@ public static void generate(PulsarService pulsar, boolean includeTopicMetrics, b topicsMap.forEach((name, topic) -> { getTopicStats(topic, topicStats, includeConsumerMetrics, pulsar.getConfiguration().isExposePreciseBacklogInPrometheus(), - pulsar.getConfiguration().isExposeSubscriptionBacklokSizeInPrometheus()); + pulsar.getConfiguration().isExposeSubscriptionBacklogSizeInPrometheus()); if (includeTopicMetrics) { topicsCount.add(1); diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java index f70d2bee10a7f..6952f8c4008ef 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java @@ -769,6 +769,10 @@ default CompletableFuture deleteAsync(String topic, boolean force) { 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, false); } From 392028bb70cdebb7a5e0bb297139e82b036f8646 Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Sun, 31 Jan 2021 20:58:02 -0800 Subject: [PATCH 5/6] Fix cmd topic npe. --- .../org/apache/pulsar/broker/admin/v2/PersistentTopics.java | 2 +- .../src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) 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 1dbcec08d9af6..fb55e823dda00 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 @@ -932,7 +932,7 @@ public TopicStats getStats( @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, @ApiParam(value = "If return precise backlog or imprecise backlog") @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog, - @ApiParam(value = "If return backlob size for each subscription, require locking on ledger so be careful " + @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); 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 07a62cc5dfaca..11cd17c5f6f73 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 @@ -542,7 +542,7 @@ private class GetStats extends CliCommand { @Override void run() throws PulsarAdminException { String topic = validateTopicName(params); - print(topics.getStats(topic, getPreciseBacklog, subscriptionBacklogSize)); + print(getTopics().getStats(topic, getPreciseBacklog, subscriptionBacklogSize)); } } From 7f0cde588a1b7adecd93803f258aee215dea194f Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Sat, 6 Feb 2021 22:16:09 -0800 Subject: [PATCH 6/6] Fix typo. --- .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 9283183fc9fd5..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 @@ -1841,7 +1841,7 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext( category = CATEGORY_METRICS, - doc = "Enable expose the backlog size fro each subscription when generating stats.\n" + + 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;