This is an automated email from the ASF dual-hosted git repository.

jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 59b019ae204 Fix consensus subscription metrics by region (#18277)
59b019ae204 is described below

commit 59b019ae204f7c9519a28e47dc449258ed49400f
Author: Caideyipi <[email protected]>
AuthorDate: Wed Jul 22 18:00:46 2026 +0800

    Fix consensus subscription metrics by region (#18277)
---
 .../consensus/ConsensusPrefetchingQueue.java       |  20 +-
 ...nsensusSubscriptionPrefetchingQueueMetrics.java | 214 +++++++++++++--------
 .../receiver/SubscriptionReceiverV1.java           |   2 +-
 .../consensus/ConsensusPrefetchingQueueTest.java   |  70 +++++++
 ...susSubscriptionPrefetchingQueueMetricsTest.java | 164 ++++++++++++++++
 .../iotdb/commons/service/metric/enums/Metric.java |   3 +
 6 files changed, 389 insertions(+), 84 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
index 19158a5b4a9..ff04adf344a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
@@ -3044,8 +3044,7 @@ public class ConsensusPrefetchingQueue {
       }
 
       // Deregister metrics after the queue is fully closed.
-      ConsensusSubscriptionPrefetchingQueueMetrics.getInstance()
-          .deregister(getPrefetchingQueueId());
+      
ConsensusSubscriptionPrefetchingQueueMetrics.getInstance().deregister(this);
 
       if (Objects.nonNull(prefetchBinding.left) && 
Objects.nonNull(prefetchBinding.right)) {
         if (!prefetchBinding.left.isShutdown()) {
@@ -3182,6 +3181,14 @@ public class ConsensusPrefetchingQueue {
     return isActive;
   }
 
+  public long getActiveStatus() {
+    return isActive ? 1L : 0L;
+  }
+
+  public long getInitializedStatus() {
+    return prefetchInitialized ? 1L : 0L;
+  }
+
   public void setActiveWriterNodeIds(final Set<Integer> activeWriterNodeIds) {
     this.runtimeActiveWriterNodeIds =
         Collections.unmodifiableSet(
@@ -3409,7 +3416,8 @@ public class ConsensusPrefetchingQueue {
         prefetchingQueue.size()
             + inFlightEvents.size()
             + pendingEntries.size()
-            + getRealtimeBufferedEntryCount();
+            + getRealtimeBufferedEntryCount()
+            + lingerBatch.getEntryCount();
     final boolean hasUnreadWalEntries = hasUnreadWalEntriesBehindCursor();
     return queuedLag + (hasUnreadWalEntries ? 1 : 0);
   }
@@ -3503,7 +3511,7 @@ public class ConsensusPrefetchingQueue {
     private long physicalTime;
     private int writerNodeId;
     private long lastLocalSeq;
-    private int entryCount;
+    private volatile int entryCount;
 
     private DeliveryBatchState() {
       reset();
@@ -3513,6 +3521,10 @@ public class ConsensusPrefetchingQueue {
       return tablets.isEmpty();
     }
 
+    private int getEntryCount() {
+      return entryCount;
+    }
+
     private void append(
         final WriterBufferedEntry entry,
         final long entryEstimatedBytes,
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/metric/ConsensusSubscriptionPrefetchingQueueMetrics.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/metric/ConsensusSubscriptionPrefetchingQueueMetrics.java
index 1172e16bd55..50ea8402118 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/metric/ConsensusSubscriptionPrefetchingQueueMetrics.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/metric/ConsensusSubscriptionPrefetchingQueueMetrics.java
@@ -44,24 +44,25 @@ public class ConsensusSubscriptionPrefetchingQueueMetrics 
implements IMetricSet
 
   private volatile AbstractMetricService metricService;
 
-  private final Map<String, ConsensusPrefetchingQueue> queueMap = new 
ConcurrentHashMap<>();
+  private final Map<QueueMetricsKey, ConsensusPrefetchingQueue> queueMap =
+      new ConcurrentHashMap<>();
 
-  private final Map<String, Rate> rateMap = new ConcurrentHashMap<>();
+  private final Map<QueueMetricsKey, Rate> rateMap = new ConcurrentHashMap<>();
 
   @Override
-  public void bindTo(final AbstractMetricService metricService) {
+  public synchronized void bindTo(final AbstractMetricService metricService) {
     this.metricService = metricService;
-    final ImmutableSet<String> ids = ImmutableSet.copyOf(queueMap.keySet());
-    for (final String id : ids) {
-      createMetrics(id);
+    final ImmutableSet<QueueMetricsKey> keys = 
ImmutableSet.copyOf(queueMap.keySet());
+    for (final QueueMetricsKey key : keys) {
+      createMetrics(key);
     }
   }
 
   @Override
-  public void unbindFrom(final AbstractMetricService metricService) {
-    final ImmutableSet<String> ids = ImmutableSet.copyOf(queueMap.keySet());
-    for (final String id : ids) {
-      deregister(id);
+  public synchronized void unbindFrom(final AbstractMetricService 
metricService) {
+    final ImmutableSet<QueueMetricsKey> keys = 
ImmutableSet.copyOf(queueMap.keySet());
+    for (final QueueMetricsKey key : keys) {
+      deregister(key);
     }
     if (!queueMap.isEmpty()) {
       LOGGER.warn(
@@ -72,21 +73,20 @@ public class ConsensusSubscriptionPrefetchingQueueMetrics 
implements IMetricSet
 
   //////////////////////////// register & deregister 
////////////////////////////
 
-  public void register(final ConsensusPrefetchingQueue queue) {
-    final String id = queue.getPrefetchingQueueId();
-    queueMap.putIfAbsent(id, queue);
-    if (Objects.nonNull(metricService)) {
-      createMetrics(id);
+  public synchronized void register(final ConsensusPrefetchingQueue queue) {
+    final QueueMetricsKey key = QueueMetricsKey.from(queue);
+    if (Objects.isNull(queueMap.putIfAbsent(key, queue)) && 
Objects.nonNull(metricService)) {
+      createMetrics(key);
     }
   }
 
-  private void createMetrics(final String id) {
-    createAutoGauge(id);
-    createRate(id);
+  private void createMetrics(final QueueMetricsKey key) {
+    createAutoGauge(key);
+    createRate(key);
   }
 
-  private void createAutoGauge(final String id) {
-    final ConsensusPrefetchingQueue queue = queueMap.get(id);
+  private void createAutoGauge(final QueueMetricsKey key) {
+    final ConsensusPrefetchingQueue queue = queueMap.get(key);
     if (Objects.isNull(queue)) {
       return;
     }
@@ -95,143 +95,199 @@ public class ConsensusSubscriptionPrefetchingQueueMetrics 
implements IMetricSet
         MetricLevel.IMPORTANT,
         queue,
         ConsensusPrefetchingQueue::getSubscriptionUncommittedEventCount,
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        key.getTags());
     // Keep the legacy metric name for dashboard compatibility, but expose 
seek generation here.
     metricService.createAutoGauge(
         Metric.SUBSCRIPTION_CURRENT_COMMIT_ID.toString(),
         MetricLevel.IMPORTANT,
         queue,
         ConsensusPrefetchingQueue::getCurrentSeekGeneration,
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        key.getTags());
+    metricService.createAutoGauge(
+        Metric.SUBSCRIPTION_CONSENSUS_SEEK_GENERATION.toString(),
+        MetricLevel.IMPORTANT,
+        queue,
+        ConsensusPrefetchingQueue::getCurrentSeekGeneration,
+        key.getTags());
     metricService.createAutoGauge(
         Metric.SUBSCRIPTION_CONSENSUS_LAG.toString(),
         MetricLevel.IMPORTANT,
         queue,
         ConsensusPrefetchingQueue::getLag,
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        key.getTags());
     metricService.createAutoGauge(
         Metric.SUBSCRIPTION_CONSENSUS_WAL_GAP.toString(),
         MetricLevel.IMPORTANT,
         queue,
         ConsensusPrefetchingQueue::getWalGapSkippedEntries,
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        key.getTags());
     metricService.createAutoGauge(
         Metric.SUBSCRIPTION_CONSENSUS_ROUTING_EPOCH_CHANGE.toString(),
         MetricLevel.IMPORTANT,
         queue,
         ConsensusPrefetchingQueue::getEpochChangeCount,
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        key.getTags());
     metricService.createAutoGauge(
         Metric.SUBSCRIPTION_CONSENSUS_WATERMARK.toString(),
         MetricLevel.IMPORTANT,
         queue,
         ConsensusPrefetchingQueue::getMaxObservedTimestamp,
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        key.getTags());
+    metricService.createAutoGauge(
+        Metric.SUBSCRIPTION_CONSENSUS_ACTIVE.toString(),
+        MetricLevel.IMPORTANT,
+        queue,
+        ConsensusPrefetchingQueue::getActiveStatus,
+        key.getTags());
+    metricService.createAutoGauge(
+        Metric.SUBSCRIPTION_CONSENSUS_INITIALIZED.toString(),
+        MetricLevel.IMPORTANT,
+        queue,
+        ConsensusPrefetchingQueue::getInitializedStatus,
+        key.getTags());
   }
 
-  private void createRate(final String id) {
-    final ConsensusPrefetchingQueue queue = queueMap.get(id);
+  private void createRate(final QueueMetricsKey key) {
+    final ConsensusPrefetchingQueue queue = queueMap.get(key);
     if (Objects.isNull(queue)) {
       return;
     }
     rateMap.put(
-        id,
+        key,
         metricService.getOrCreateRate(
-            Metric.SUBSCRIPTION_EVENT_TRANSFER.toString(),
-            MetricLevel.IMPORTANT,
-            Tag.NAME.toString(),
-            queue.getPrefetchingQueueId()));
+            Metric.SUBSCRIPTION_EVENT_TRANSFER.toString(), 
MetricLevel.IMPORTANT, key.getTags()));
   }
 
-  public void deregister(final String id) {
-    if (!queueMap.containsKey(id)) {
+  public synchronized void deregister(final ConsensusPrefetchingQueue queue) {
+    final QueueMetricsKey key = QueueMetricsKey.from(queue);
+    if (queueMap.get(key) != queue) {
       LOGGER.warn(
           DataNodePipeMessages
               
.PIPE_LOG_FAILED_TO_DEREGISTER_CONSENSUS_SUBSCRIPTION_PREFETCHING_8B180091,
-          id);
+          key);
+      return;
+    }
+    deregister(key);
+  }
+
+  private void deregister(final QueueMetricsKey key) {
+    if (!queueMap.containsKey(key)) {
+      LOGGER.warn(
+          DataNodePipeMessages
+              
.PIPE_LOG_FAILED_TO_DEREGISTER_CONSENSUS_SUBSCRIPTION_PREFETCHING_8B180091,
+          key);
       return;
     }
     if (Objects.nonNull(metricService)) {
-      removeMetrics(id);
+      removeMetrics(key);
     }
-    queueMap.remove(id);
+    queueMap.remove(key);
   }
 
-  private void removeMetrics(final String id) {
-    removeAutoGauge(id);
-    removeRate(id);
+  private void removeMetrics(final QueueMetricsKey key) {
+    removeAutoGauge(key);
+    removeRate(key);
   }
 
-  private void removeAutoGauge(final String id) {
-    final ConsensusPrefetchingQueue queue = queueMap.get(id);
-    if (Objects.isNull(queue)) {
+  private void removeAutoGauge(final QueueMetricsKey key) {
+    if (!queueMap.containsKey(key)) {
       return;
     }
     metricService.remove(
         MetricType.AUTO_GAUGE,
         Metric.SUBSCRIPTION_UNCOMMITTED_EVENT_COUNT.toString(),
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        key.getTags());
     metricService.remove(
-        MetricType.AUTO_GAUGE,
-        Metric.SUBSCRIPTION_CURRENT_COMMIT_ID.toString(),
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        MetricType.AUTO_GAUGE, 
Metric.SUBSCRIPTION_CURRENT_COMMIT_ID.toString(), key.getTags());
     metricService.remove(
         MetricType.AUTO_GAUGE,
-        Metric.SUBSCRIPTION_CONSENSUS_LAG.toString(),
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        Metric.SUBSCRIPTION_CONSENSUS_SEEK_GENERATION.toString(),
+        key.getTags());
     metricService.remove(
-        MetricType.AUTO_GAUGE,
-        Metric.SUBSCRIPTION_CONSENSUS_WAL_GAP.toString(),
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        MetricType.AUTO_GAUGE, Metric.SUBSCRIPTION_CONSENSUS_LAG.toString(), 
key.getTags());
+    metricService.remove(
+        MetricType.AUTO_GAUGE, 
Metric.SUBSCRIPTION_CONSENSUS_WAL_GAP.toString(), key.getTags());
     metricService.remove(
         MetricType.AUTO_GAUGE,
         Metric.SUBSCRIPTION_CONSENSUS_ROUTING_EPOCH_CHANGE.toString(),
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        key.getTags());
     metricService.remove(
-        MetricType.AUTO_GAUGE,
-        Metric.SUBSCRIPTION_CONSENSUS_WATERMARK.toString(),
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        MetricType.AUTO_GAUGE, 
Metric.SUBSCRIPTION_CONSENSUS_WATERMARK.toString(), key.getTags());
+    metricService.remove(
+        MetricType.AUTO_GAUGE, 
Metric.SUBSCRIPTION_CONSENSUS_ACTIVE.toString(), key.getTags());
+    metricService.remove(
+        MetricType.AUTO_GAUGE, 
Metric.SUBSCRIPTION_CONSENSUS_INITIALIZED.toString(), key.getTags());
   }
 
-  private void removeRate(final String id) {
-    final ConsensusPrefetchingQueue queue = queueMap.get(id);
-    if (Objects.isNull(queue)) {
+  private void removeRate(final QueueMetricsKey key) {
+    if (!queueMap.containsKey(key)) {
       return;
     }
     metricService.remove(
-        MetricType.RATE,
-        Metric.SUBSCRIPTION_EVENT_TRANSFER.toString(),
-        Tag.NAME.toString(),
-        queue.getPrefetchingQueueId());
+        MetricType.RATE, Metric.SUBSCRIPTION_EVENT_TRANSFER.toString(), 
key.getTags());
+    rateMap.remove(key);
   }
 
-  public void mark(final String id, final long size) {
+  public void mark(final String id, final String regionId, final long size) {
     if (Objects.isNull(metricService)) {
       return;
     }
-    final Rate rate = rateMap.get(id);
+    final QueueMetricsKey key = new QueueMetricsKey(id, regionId);
+    final Rate rate = rateMap.get(key);
     if (rate == null) {
       LOGGER.warn(
           DataNodePipeMessages
               
.PIPE_LOG_FAILED_TO_MARK_TRANSFER_EVENT_RATE_CONSENSUSPREFETCHINGQUEUE_FE9B91C3,
-          id);
+          key);
       return;
     }
     rate.mark(size);
   }
 
+  private static final class QueueMetricsKey {
+
+    private final String queueId;
+    private final String regionId;
+
+    private QueueMetricsKey(final String queueId, final String regionId) {
+      this.queueId = queueId;
+      this.regionId = regionId;
+    }
+
+    private static QueueMetricsKey from(final ConsensusPrefetchingQueue queue) 
{
+      return new QueueMetricsKey(
+          queue.getPrefetchingQueueId(), 
queue.getConsensusGroupId().toString());
+    }
+
+    private String[] getTags() {
+      return new String[] {
+        Tag.NAME.toString(), queueId, Tag.REGION.toString(), regionId,
+      };
+    }
+
+    @Override
+    public boolean equals(final Object obj) {
+      if (this == obj) {
+        return true;
+      }
+      if (!(obj instanceof QueueMetricsKey)) {
+        return false;
+      }
+      final QueueMetricsKey that = (QueueMetricsKey) obj;
+      return Objects.equals(queueId, that.queueId) && Objects.equals(regionId, 
that.regionId);
+    }
+
+    @Override
+    public int hashCode() {
+      return Objects.hash(queueId, regionId);
+    }
+
+    @Override
+    public String toString() {
+      return queueId + "/" + regionId;
+    }
+  }
+
   //////////////////////////// singleton ////////////////////////////
 
   private static class Holder {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
index 696e2aa9b39..c7c791ff128 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
@@ -660,7 +660,7 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
                     if 
(ConsensusSubscriptionSetupHandler.isConsensusBasedTopic(
                         commitContext.getTopicName())) {
                       
ConsensusSubscriptionPrefetchingQueueMetrics.getInstance()
-                          .mark(queueId, size);
+                          .mark(queueId, commitContext.getRegionId(), size);
                     } else {
                       
SubscriptionPrefetchingQueueMetrics.getInstance().mark(queueId, size);
                     }
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
index 41ad9143643..4902554dc90 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
@@ -19,6 +19,7 @@
 
 package org.apache.iotdb.db.subscription.broker.consensus;
 
+import org.apache.iotdb.commons.conf.CommonDescriptor;
 import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
 import org.apache.iotdb.consensus.iot.IoTConsensusServerImpl;
@@ -74,6 +75,75 @@ public class ConsensusPrefetchingQueueTest {
 
   @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();
 
+  @Test
+  public void testLagIncludesLingeringBatchUntilCommitted() throws Exception {
+    final String originalSystemDir = 
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
+    final int originalBatchMaxDelay =
+        
CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxDelayInMs();
+    final File systemDir = temporaryFolder.newFolder("lagWithLingeringBatch");
+    ConsensusPrefetchingQueue queue = null;
+    try {
+      
CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxDelayInMs(60_000);
+      final DataRegionId regionId = new DataRegionId(9);
+      final FakeConsensusReqReader reader = new FakeConsensusReqReader();
+      reader.currentSearchIndex = 1L;
+      final IoTConsensusServerImpl serverImpl = 
mock(IoTConsensusServerImpl.class);
+      when(serverImpl.getConsensusReqReader()).thenReturn(reader);
+      when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new 
WriterSafeFrontierTracker());
+      final ConsensusLogToTabletConverter converter = 
mock(ConsensusLogToTabletConverter.class);
+      
when(converter.convert(any())).thenReturn(Collections.singletonList(createTablet()));
+      when(converter.getDatabaseName()).thenReturn("db");
+      queue =
+          new ConsensusPrefetchingQueue(
+              "consumerGroup",
+              "topic",
+              TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
+              regionId,
+              serverImpl,
+              new SubscriptionWalRetentionPolicy(
+                  "topic",
+                  SubscriptionWalRetentionPolicy.UNBOUNDED,
+                  SubscriptionWalRetentionPolicy.UNBOUNDED),
+              converter,
+              newCommitManager(systemDir),
+              new RegionProgress(Collections.emptyMap()),
+              1L,
+              1L,
+              true);
+      final IndexedConsensusRequest request =
+          new IndexedConsensusRequest(
+                  1L, 
Collections.singletonList(StatementTestUtils.genInsertRowNode(1)))
+              .setPhysicalTime(1000L)
+              .setNodeId(7);
+
+      assertNull(queue.poll("consumer"));
+      pendingEntries(queue).offer(request);
+      queue.drivePrefetchOnce();
+
+      assertEquals(0, queue.getPrefetchedEventCount());
+      assertEquals(1L, queue.getLag());
+
+      
CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxDelayInMs(0);
+      queue.drivePrefetchOnce();
+      assertEquals(1, queue.getPrefetchedEventCount());
+      assertEquals(1L, queue.getLag());
+
+      final SubscriptionEvent event = queue.poll("consumer");
+      assertNotNull(event);
+      assertEquals(1L, queue.getLag());
+      assertTrue(queue.ack("consumer", event.getCommitContext()));
+      assertEquals(0L, queue.getLag());
+    } finally {
+      if (queue != null) {
+        queue.close();
+      }
+      CommonDescriptor.getInstance()
+          .getConfig()
+          .setSubscriptionConsensusBatchMaxDelayInMs(originalBatchMaxDelay);
+      
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
+    }
+  }
+
   @Test
   public void testFilteredEmptyEntryAdvancesProgressWithoutEvent() throws 
Exception {
     final String originalSystemDir = 
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/metric/ConsensusSubscriptionPrefetchingQueueMetricsTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/metric/ConsensusSubscriptionPrefetchingQueueMetricsTest.java
new file mode 100644
index 00000000000..f6bae43f238
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/metric/ConsensusSubscriptionPrefetchingQueueMetricsTest.java
@@ -0,0 +1,164 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.subscription.metric;
+
+import org.apache.iotdb.commons.consensus.DataRegionId;
+import org.apache.iotdb.commons.service.metric.enums.Metric;
+import org.apache.iotdb.commons.service.metric.enums.Tag;
+import 
org.apache.iotdb.db.subscription.broker.consensus.ConsensusPrefetchingQueue;
+import org.apache.iotdb.metrics.AbstractMetricService;
+import org.apache.iotdb.metrics.type.Rate;
+import org.apache.iotdb.metrics.utils.MetricLevel;
+import org.apache.iotdb.metrics.utils.MetricType;
+
+import org.junit.Test;
+
+import java.lang.reflect.Field;
+import java.util.Map;
+
+import static org.junit.Assert.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class ConsensusSubscriptionPrefetchingQueueMetricsTest {
+
+  @Test
+  public void testMetricsAreIsolatedByRegion() throws Exception {
+    final String queueId = "consumer_group_topic";
+    final DataRegionId firstRegionId = new DataRegionId(1);
+    final DataRegionId secondRegionId = new DataRegionId(2);
+    final ConsensusPrefetchingQueue firstQueue = 
mock(ConsensusPrefetchingQueue.class);
+    final ConsensusPrefetchingQueue secondQueue = 
mock(ConsensusPrefetchingQueue.class);
+    final ConsensusPrefetchingQueue staleFirstQueue = 
mock(ConsensusPrefetchingQueue.class);
+    final AbstractMetricService metricService = 
mock(AbstractMetricService.class);
+    final Rate firstRate = mock(Rate.class);
+    final Rate secondRate = mock(Rate.class);
+
+    when(firstQueue.getPrefetchingQueueId()).thenReturn(queueId);
+    when(firstQueue.getConsensusGroupId()).thenReturn(firstRegionId);
+    when(secondQueue.getPrefetchingQueueId()).thenReturn(queueId);
+    when(secondQueue.getConsensusGroupId()).thenReturn(secondRegionId);
+    when(staleFirstQueue.getPrefetchingQueueId()).thenReturn(queueId);
+    when(staleFirstQueue.getConsensusGroupId()).thenReturn(firstRegionId);
+    when(metricService.getOrCreateRate(
+            Metric.SUBSCRIPTION_EVENT_TRANSFER.toString(),
+            MetricLevel.IMPORTANT,
+            Tag.NAME.toString(),
+            queueId,
+            Tag.REGION.toString(),
+            firstRegionId.toString()))
+        .thenReturn(firstRate);
+    when(metricService.getOrCreateRate(
+            Metric.SUBSCRIPTION_EVENT_TRANSFER.toString(),
+            MetricLevel.IMPORTANT,
+            Tag.NAME.toString(),
+            queueId,
+            Tag.REGION.toString(),
+            secondRegionId.toString()))
+        .thenReturn(secondRate);
+
+    final ConsensusSubscriptionPrefetchingQueueMetrics metrics =
+        ConsensusSubscriptionPrefetchingQueueMetrics.getInstance();
+    final Field metricServiceField =
+        
ConsensusSubscriptionPrefetchingQueueMetrics.class.getDeclaredField("metricService");
+    final Field queueMapField =
+        
ConsensusSubscriptionPrefetchingQueueMetrics.class.getDeclaredField("queueMap");
+    final Field rateMapField =
+        
ConsensusSubscriptionPrefetchingQueueMetrics.class.getDeclaredField("rateMap");
+    metricServiceField.setAccessible(true);
+    queueMapField.setAccessible(true);
+    rateMapField.setAccessible(true);
+    final Map<?, ?> queueMap = (Map<?, ?>) queueMapField.get(metrics);
+    final Map<?, ?> rateMap = (Map<?, ?>) rateMapField.get(metrics);
+    queueMap.clear();
+    rateMap.clear();
+    metricServiceField.set(metrics, null);
+
+    try {
+      metrics.bindTo(metricService);
+      metrics.register(firstQueue);
+      metrics.register(secondQueue);
+
+      assertEquals(2, queueMap.size());
+      assertEquals(2, rateMap.size());
+      verify(metricService)
+          .createAutoGauge(
+              eq(Metric.SUBSCRIPTION_CONSENSUS_LAG.toString()),
+              eq(MetricLevel.IMPORTANT),
+              eq(firstQueue),
+              any(),
+              eq(Tag.NAME.toString()),
+              eq(queueId),
+              eq(Tag.REGION.toString()),
+              eq(firstRegionId.toString()));
+      verify(metricService)
+          .createAutoGauge(
+              eq(Metric.SUBSCRIPTION_CONSENSUS_LAG.toString()),
+              eq(MetricLevel.IMPORTANT),
+              eq(secondQueue),
+              any(),
+              eq(Tag.NAME.toString()),
+              eq(queueId),
+              eq(Tag.REGION.toString()),
+              eq(secondRegionId.toString()));
+
+      metrics.mark(queueId, firstRegionId.toString(), 11L);
+      metrics.mark(queueId, secondRegionId.toString(), 22L);
+      verify(firstRate).mark(11L);
+      verify(secondRate).mark(22L);
+
+      metrics.register(staleFirstQueue);
+      metrics.deregister(staleFirstQueue);
+      assertEquals(2, queueMap.size());
+      assertEquals(2, rateMap.size());
+
+      metrics.deregister(firstQueue);
+      assertEquals(1, queueMap.size());
+      assertEquals(1, rateMap.size());
+      verify(metricService)
+          .remove(
+              MetricType.AUTO_GAUGE,
+              Metric.SUBSCRIPTION_CONSENSUS_LAG.toString(),
+              Tag.NAME.toString(),
+              queueId,
+              Tag.REGION.toString(),
+              firstRegionId.toString());
+      verify(metricService, never())
+          .remove(
+              MetricType.AUTO_GAUGE,
+              Metric.SUBSCRIPTION_CONSENSUS_LAG.toString(),
+              Tag.NAME.toString(),
+              queueId,
+              Tag.REGION.toString(),
+              secondRegionId.toString());
+
+      metrics.mark(queueId, secondRegionId.toString(), 33L);
+      verify(secondRate).mark(33L);
+    } finally {
+      queueMap.clear();
+      rateMap.clear();
+      metricServiceField.set(metrics, null);
+    }
+  }
+}
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/enums/Metric.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/enums/Metric.java
index f5a76e6d66a..238af422299 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/enums/Metric.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/enums/Metric.java
@@ -212,6 +212,9 @@ public enum Metric {
   SUBSCRIPTION_CONSENSUS_WAL_GAP("subscription_consensus_wal_gap"),
   
SUBSCRIPTION_CONSENSUS_ROUTING_EPOCH_CHANGE("subscription_consensus_routing_epoch_change"),
   SUBSCRIPTION_CONSENSUS_WATERMARK("subscription_consensus_watermark"),
+  
SUBSCRIPTION_CONSENSUS_SEEK_GENERATION("subscription_consensus_seek_generation"),
+  SUBSCRIPTION_CONSENSUS_ACTIVE("subscription_consensus_active"),
+  SUBSCRIPTION_CONSENSUS_INITIALIZED("subscription_consensus_initialized"),
   // load related
   ACTIVE_LOADING_FILES_NUMBER("active_loading_files_number"),
   ACTIVE_LOADING_FILES_SIZE("active_loading_files_size"),

Reply via email to