diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/BrokerOperabilityMetrics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/BrokerOperabilityMetrics.java index 400dbd3335a2a..09b313aa749a7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/BrokerOperabilityMetrics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/BrokerOperabilityMetrics.java @@ -33,6 +33,7 @@ public class BrokerOperabilityMetrics { private static final Counter TOPIC_LOAD_FAILED = Counter.build("topic_load_failed", "-").register(); private final List metricsList; private final String localCluster; + private final DimensionStats oldTopicLoadStats; private final DimensionStats topicLoadStats; private final String brokerName; private final LongAdder connectionTotalCreatedCount; @@ -44,7 +45,8 @@ public class BrokerOperabilityMetrics { public BrokerOperabilityMetrics(String localCluster, String brokerName) { this.metricsList = new ArrayList<>(); this.localCluster = localCluster; - this.topicLoadStats = new DimensionStats("topic_load_times", 60); + this.oldTopicLoadStats = new DimensionStats("topic_load_times", 60); + this.topicLoadStats = new DimensionStats("pulsar_topic_load_times", 60); this.brokerName = brokerName; this.connectionTotalCreatedCount = new LongAdder(); this.connectionCreateSuccessCount = new LongAdder(); @@ -59,6 +61,7 @@ public List getMetrics() { } private void generate() { + metricsList.add(getOldTopicLoadMetrics()); metricsList.add(getTopicLoadMetrics()); metricsList.add(getConnectionMetrics()); } @@ -85,8 +88,13 @@ Map getDimensionMap(String metricsName) { return dimensionMap; } + Metrics getOldTopicLoadMetrics() { + Metrics metrics = getDimensionMetrics("topic_load_times", "topic_load", oldTopicLoadStats); + return metrics; + } + Metrics getTopicLoadMetrics() { - Metrics metrics = getDimensionMetrics("topic_load_times", "topic_load", topicLoadStats); + Metrics metrics = getDimensionMetrics("pulsar_topic_load_times", "topic_load", topicLoadStats); metrics.put("brk_topic_load_failed_count", TOPIC_LOAD_FAILED.get()); return metrics; } @@ -109,10 +117,12 @@ Metrics getDimensionMetrics(String metricsName, String dimensionName, DimensionS public void reset() { metricsList.clear(); + oldTopicLoadStats.reset(); topicLoadStats.reset(); } public void recordTopicLoadTimeValue(long topicLoadLatencyMs) { + oldTopicLoadStats.recordDimensionTimeValue(topicLoadLatencyMs, TimeUnit.MILLISECONDS); topicLoadStats.recordDimensionTimeValue(topicLoadLatencyMs, TimeUnit.MILLISECONDS); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/DimensionStats.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/DimensionStats.java index 604265b554050..1b6f981ca4e21 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/DimensionStats.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/DimensionStats.java @@ -64,7 +64,7 @@ public DimensionStats(String name, long updateDurationInSec) { defaultRegistry.register(summary); } catch (IllegalArgumentException ie) { // it only happens in test-cases when try to register summary multiple times in registry - log.warn("{} is already registred {}", name, ie.getMessage()); + log.warn("{} is already registered {}", name, ie.getMessage()); } } } 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 6cb7378330f09..c4e41074a1ab3 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 @@ -246,6 +246,14 @@ public void testMetricsTopicCount() throws Exception { assertEquals(item.value, 3.0); } }); + Collection topicLoadTimesMetrics = metrics.get("topic_load_times"); + Collection topicLoadTimesCountMetrics = metrics.get("topic_load_times_count"); + assertEquals(topicLoadTimesMetrics.size(), 6); + assertEquals(topicLoadTimesCountMetrics.size(), 1); + Collection pulsarTopicLoadTimesMetrics = metrics.get("pulsar_topic_load_times"); + Collection pulsarTopicLoadTimesCountMetrics = metrics.get("pulsar_topic_load_times_count"); + assertEquals(pulsarTopicLoadTimesMetrics.size(), 6); + assertEquals(pulsarTopicLoadTimesCountMetrics.size(), 1); } @Test