From 4dfb6d4c451b12ba815ba1c0ef07eba4ea65a78a Mon Sep 17 00:00:00 2001 From: shawyeok Date: Sat, 24 Dec 2022 21:09:17 +0800 Subject: [PATCH] Fix configuration exposeSubscriptionBacklogSizeInPrometheus not working --- .../AggregatedSubscriptionStats.java | 2 + .../prometheus/NamespaceStatsAggregator.java | 1 + .../broker/stats/prometheus/TopicStats.java | 4 ++ .../broker/stats/PrometheusMetricsTest.java | 45 +++++++++++++++++++ .../data/stats/SubscriptionStatsImpl.java | 8 ++-- 5 files changed, 57 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/AggregatedSubscriptionStats.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/AggregatedSubscriptionStats.java index 9b37f7d48c8e5..fc41ad4853816 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/AggregatedSubscriptionStats.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/AggregatedSubscriptionStats.java @@ -28,6 +28,8 @@ public class AggregatedSubscriptionStats { public long msgBacklogNoDelayed; + public long msgBacklogSize; + public boolean blockedSubscriptionOnUnackedMsgs; public double msgRateRedeliver; 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 26aae969aab42..529666f6e84f4 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 @@ -119,6 +119,7 @@ private static void aggregateTopicStats(TopicStats stats, SubscriptionStatsImpl subsStats.msgOutCounter = subscriptionStats.msgOutCounter; subsStats.msgBacklog = subscriptionStats.msgBacklog; subsStats.msgDelayed = subscriptionStats.msgDelayed; + subsStats.msgBacklogSize = subscriptionStats.backlogSize; subsStats.msgRateExpired = subscriptionStats.msgRateExpired; subsStats.totalMsgExpired = subscriptionStats.totalMsgExpired; subsStats.msgBacklogNoDelayed = subsStats.msgBacklog - subsStats.msgDelayed; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/TopicStats.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/TopicStats.java index 754722c98be77..844d866253713 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/TopicStats.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/TopicStats.java @@ -254,6 +254,10 @@ public static void printTopicStats(PrometheusMetricStreams stream, TopicStats st stats.subscriptionStats.forEach((sub, subsStats) -> { writeSubscriptionMetric(stream, "pulsar_subscription_back_log", subsStats.msgBacklog, cluster, namespace, topic, sub, splitTopicAndPartitionIndexLabel); + if (subsStats.msgBacklogSize >= 0) { + writeSubscriptionMetric(stream, "pulsar_subscription_back_log_size", subsStats.msgBacklogSize, + cluster, namespace, topic, sub, splitTopicAndPartitionIndexLabel); + } writeSubscriptionMetric(stream, "pulsar_subscription_back_log_no_delayed", subsStats.msgBacklogNoDelayed, cluster, namespace, topic, sub, splitTopicAndPartitionIndexLabel); writeSubscriptionMetric(stream, "pulsar_subscription_delayed", diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java index db59942e15ce5..44e722af24797 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java @@ -20,6 +20,7 @@ import static com.google.common.base.Preconditions.checkArgument; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; @@ -50,6 +51,7 @@ import java.util.Set; import java.util.TreeMap; import java.util.UUID; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.regex.Matcher; @@ -72,12 +74,17 @@ import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.stats.prometheus.PrometheusMetricsGenerator; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.common.naming.NamespaceName; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.compaction.Compactor; import org.awaitility.Awaitility; import org.mockito.Mockito; @@ -606,6 +613,7 @@ public void testNonPersistentSubMetrics() throws Exception { String metricsStr = statsOut.toString(); Multimap metrics = parseMetrics(metricsStr); assertTrue(metrics.containsKey("pulsar_subscription_back_log")); + assertTrue(metrics.containsKey("pulsar_subscription_back_log_size")); assertTrue(metrics.containsKey("pulsar_subscription_back_log_no_delayed")); assertTrue(metrics.containsKey("pulsar_subscription_msg_throughput_out")); assertTrue(metrics.containsKey("pulsar_throughput_out")); @@ -1689,6 +1697,43 @@ public void testMetricsGroupedByTypeDefinitions() throws Exception { p2.close(); } + @Test + public void testSubscriptionBacklogSizeMetric() throws PulsarAdminException, IOException { + NamespaceName namespaceName = NamespaceName.get("prop/ns-abc"); + TopicName topicName = TopicName.get("persistent", namespaceName, UUID.randomUUID().toString()); + admin.topics().createNonPartitionedTopic(topicName.toString()); + String subscriptionName = "sub"; + admin.topics().createSubscription(topicName.toString(), subscriptionName, MessageId.latest); + @Cleanup + Producer producer = pulsarClient.newProducer() + .topic(topicName.toString()) + .enableBatching(false) + .create(); + List> futures = new ArrayList<>(); + for (int i = 0; i < 10; i++) { + futures.add(producer.newMessage().value(new byte[1024]).sendAsync()); + } + FutureUtil.waitForAll(futures).join(); + Multimap metrics; + pulsar.getConfiguration().setExposeSubscriptionBacklogSizeInPrometheus(false); + ByteArrayOutputStream statsOut1 = new ByteArrayOutputStream(); + PrometheusMetricsGenerator.generate(pulsar, true, false, false, statsOut1); + metrics = parseMetrics(statsOut1.toString()); + assertFalse(metrics.containsKey("pulsar_subscription_back_log_size")); + + pulsar.getConfiguration().setExposeSubscriptionBacklogSizeInPrometheus(true); + ByteArrayOutputStream statsOut2 = new ByteArrayOutputStream(); + PrometheusMetricsGenerator.generate(pulsar, true, false, false, statsOut2); + metrics = parseMetrics(statsOut2.toString()); + assertTrue(metrics.containsKey("pulsar_subscription_back_log_size")); + Optional metric = metrics.get("pulsar_subscription_back_log_size").stream() + .filter(item -> item.tags.get("topic").equals(topicName.toString()) + && item.tags.get("subscription").equals(subscriptionName)) + .findAny(); + assertTrue(metric.isPresent()); + assertTrue(metric.get().value > 0); + } + /** * Hacky parsing of Prometheus text format. Should be good enough for unit tests */ diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/SubscriptionStatsImpl.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/SubscriptionStatsImpl.java index 8eb77e981a883..ad8e2c8276e90 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/SubscriptionStatsImpl.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/SubscriptionStatsImpl.java @@ -59,7 +59,7 @@ public class SubscriptionStatsImpl implements SubscriptionStats { public long msgBacklog; /** Size of backlog in byte. **/ - public long backlogSize; + public long backlogSize = -1; /** Get the publish time of the earliest message in the backlog. */ public long earliestMsgPublishTimeInBacklog; @@ -158,7 +158,7 @@ public void reset() { msgOutCounter = 0; msgRateRedeliver = 0; msgBacklog = 0; - backlogSize = 0; + backlogSize = -1; msgBacklogNoDelayed = 0; unackedMessages = 0; msgRateExpired = 0; @@ -187,7 +187,9 @@ public SubscriptionStatsImpl add(SubscriptionStatsImpl stats) { this.msgOutCounter += stats.msgOutCounter; this.msgRateRedeliver += stats.msgRateRedeliver; this.msgBacklog += stats.msgBacklog; - this.backlogSize += stats.backlogSize; + if (stats.backlogSize >= 0) { + this.backlogSize = Math.max(this.backlogSize, 0) + stats.backlogSize; + } this.msgBacklogNoDelayed += stats.msgBacklogNoDelayed; this.msgDelayed += stats.msgDelayed; this.unackedMessages += stats.unackedMessages;