This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git
The following commit(s) were added to refs/heads/develop by this push:
new 577b89f2cd [ISSUE #10644] Add LMQ number gauge metric to Broker
observability (#10645)
577b89f2cd is described below
commit 577b89f2cdddf0d42cfd3c1b6effdac0cd0e467c
Author: Quan <[email protected]>
AuthorDate: Wed Jul 22 11:43:58 2026 +0800
[ISSUE #10644] Add LMQ number gauge metric to Broker observability (#10645)
- Add GAUGE_LMQ_NUM constant in BrokerMetricsConstant
- Register ObservableLongGauge in BrokerMetricsManager reading from
messageStore.getQueueStore().getLmqNum()
- Expose rocketmq_lmq_number metric for monitoring LMQ resource usage
---
.../apache/rocketmq/broker/metrics/BrokerMetricsConstant.java | 1 +
.../org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java | 9 +++++++++
2 files changed, 10 insertions(+)
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsConstant.java
b/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsConstant.java
index e87ce9ad02..d39416e4e6 100644
---
a/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsConstant.java
+++
b/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsConstant.java
@@ -25,6 +25,7 @@ public class BrokerMetricsConstant {
public static final String GAUGE_BROKER_PERMISSION =
"rocketmq_broker_permission";
public static final String GAUGE_TOPIC_NUM = "rocketmq_topic_number";
public static final String GAUGE_CONSUMER_GROUP_NUM =
"rocketmq_consumer_group_number";
+ public static final String GAUGE_LMQ_NUM = "rocketmq_lmq_number";
public static final String COUNTER_MESSAGES_IN_TOTAL =
"rocketmq_messages_in_total";
public static final String COUNTER_MESSAGES_OUT_TOTAL =
"rocketmq_messages_out_total";
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java
b/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java
index eda80dfd44..299b712b37 100644
---
a/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java
+++
b/broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java
@@ -93,6 +93,7 @@ import static
org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.GAUGE_CON
import static
org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.GAUGE_CONSUMER_QUEUEING_LATENCY;
import static
org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.GAUGE_CONSUMER_READY_MESSAGES;
import static
org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.GAUGE_HALF_MESSAGES;
+import static
org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.GAUGE_LMQ_NUM;
import static
org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.GAUGE_PROCESSOR_WATERMARK;
import static
org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.GAUGE_PRODUCER_CONNECTIONS;
import static
org.apache.rocketmq.broker.metrics.BrokerMetricsConstant.HISTOGRAM_FINISH_MSG_LATENCY;
@@ -138,6 +139,7 @@ public class BrokerMetricsManager {
private ObservableLongGauge brokerPermission = new
NopObservableLongGauge();
private ObservableLongGauge topicNum = new NopObservableLongGauge();
private ObservableLongGauge consumerGroupNum = new
NopObservableLongGauge();
+ private ObservableLongGauge lmqNum = new NopObservableLongGauge();
// request metrics
private LongCounter messagesInTotal = new NopLongCounter();
@@ -602,6 +604,13 @@ public class BrokerMetricsManager {
.setDescription("Active subscription group number")
.ofLongs()
.buildWithCallback(measurement ->
measurement.record(brokerController.getSubscriptionGroupManager().getSubscriptionGroupTable().size(),
newAttributesBuilder().build()));
+
+ if (messageStore.getMessageStoreConfig().isEnableLmq()) {
+ lmqNum = brokerMeter.gaugeBuilder(GAUGE_LMQ_NUM)
+ .setDescription("Current LMQ number")
+ .ofLongs()
+ .buildWithCallback(measurement ->
measurement.record(messageStore.getQueueStore().getLmqNum(),
newAttributesBuilder().build()));
+ }
}
private void initRequestMetrics() {