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

tanxinyu 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 f077a18567f [IOTDB-6281] Enhance Procedure metrics (#11798)
f077a18567f is described below

commit f077a18567feffc209b4974844319297c7acec99
Author: Xiangpeng Hu <[email protected]>
AuthorDate: Thu Jan 4 15:53:50 2024 +0800

    [IOTDB-6281] Enhance Procedure metrics (#11798)
---
 .../iotdb/confignode/manager/ConfigManager.java    |   2 +
 .../iotdb/confignode/manager/ProcedureManager.java |  16 ++
 .../iotdb/confignode/procedure/Procedure.java      |  46 +++++-
 .../confignode/procedure/ProcedureExecutor.java    |  19 ++-
 .../confignode/procedure/ProcedureMetrics.java     | 184 +++++++++++++++++++++
 .../impl/statemachine/StateMachineProcedure.java   |   2 +-
 .../procedure/TestProcedureExecutor.java           |   2 +-
 .../procedure/entity/StuckProcedure.java           |   4 +-
 .../iotdb/commons/service/metric/enums/Metric.java |   6 +
 9 files changed, 266 insertions(+), 15 deletions(-)

diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
index 32685a6ad7f..69c8729e847 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
@@ -1485,12 +1485,14 @@ public class ConfigManager implements IManager {
   public void addMetrics() {
     MetricService.getInstance().addMetricSet(new 
NodeMetrics(getNodeManager()));
     MetricService.getInstance().addMetricSet(new PartitionMetrics(this));
+    getProcedureManager().addMetrics();
   }
 
   @Override
   public void removeMetrics() {
     MetricService.getInstance().removeMetricSet(new 
NodeMetrics(getNodeManager()));
     MetricService.getInstance().removeMetricSet(new PartitionMetrics(this));
+    getProcedureManager().removeMetrics();
   }
 
   @Override
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
index 11c01a2eedb..965f89b69c7 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
@@ -34,6 +34,7 @@ import org.apache.iotdb.commons.path.PathDeserializeUtil;
 import org.apache.iotdb.commons.path.PathPatternTree;
 import org.apache.iotdb.commons.pipe.plugin.meta.PipePluginMeta;
 import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
+import org.apache.iotdb.commons.service.metric.MetricService;
 import org.apache.iotdb.commons.trigger.TriggerInformation;
 import org.apache.iotdb.commons.utils.StatusUtils;
 import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
@@ -47,6 +48,7 @@ import 
org.apache.iotdb.confignode.manager.partition.PartitionManager;
 import org.apache.iotdb.confignode.persistence.ProcedureInfo;
 import org.apache.iotdb.confignode.procedure.Procedure;
 import org.apache.iotdb.confignode.procedure.ProcedureExecutor;
+import org.apache.iotdb.confignode.procedure.ProcedureMetrics;
 import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
 import org.apache.iotdb.confignode.procedure.impl.cq.CreateCQProcedure;
 import org.apache.iotdb.confignode.procedure.impl.node.AddConfigNodeProcedure;
@@ -129,6 +131,7 @@ public class ProcedureManager {
   private ConfigNodeProcedureEnv env;
 
   private final long planSizeLimit;
+  private ProcedureMetrics procedureMetrics;
 
   public ProcedureManager(ConfigManager configManager, ProcedureInfo 
procedureInfo) {
     this.configManager = configManager;
@@ -141,6 +144,7 @@ public class ProcedureManager {
                 .getConf()
                 .getConfigNodeRatisConsensusLogAppenderBufferSize()
             - IoTDBConstant.RAFT_LOG_BASIC_SIZE;
+    this.procedureMetrics = new ProcedureMetrics(this);
   }
 
   public void shiftExecutor(boolean running) {
@@ -1027,4 +1031,16 @@ public class ProcedureManager {
               }
             });
   }
+
+  public void addMetrics() {
+    MetricService.getInstance().addMetricSet(this.procedureMetrics);
+  }
+
+  public void removeMetrics() {
+    MetricService.getInstance().removeMetricSet(this.procedureMetrics);
+  }
+
+  public ProcedureMetrics getProcedureMetrics() {
+    return procedureMetrics;
+  }
 }
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/Procedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/Procedure.java
index 1974f6d8d5b..0b90e33dad1 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/Procedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/Procedure.java
@@ -19,6 +19,7 @@
 
 package org.apache.iotdb.confignode.procedure;
 
+import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
 import 
org.apache.iotdb.confignode.procedure.exception.ProcedureAbortedException;
 import org.apache.iotdb.confignode.procedure.exception.ProcedureException;
 import 
org.apache.iotdb.confignode.procedure.exception.ProcedureSuspendedException;
@@ -508,12 +509,6 @@ public abstract class Procedure<Env> implements 
Comparable<Procedure<Env>> {
     return sb.toString();
   }
 
-  protected String toStringClass() {
-    StringBuilder sb = new StringBuilder();
-    toStringClassDetails(sb);
-    return sb.toString();
-  }
-
   /**
    * Called from {@link #toString()} when interpolating {@link Procedure} 
State. Allows decorating
    * generic Procedure State with Procedure particulars.
@@ -556,8 +551,8 @@ public abstract class Procedure<Env> implements 
Comparable<Procedure<Env>> {
     return rootProcId;
   }
 
-  public String getProcName() {
-    return toStringClass();
+  public String getProcType() {
+    return getClass().getSimpleName();
   }
 
   public long getSubmittedTime() {
@@ -900,4 +895,39 @@ public abstract class Procedure<Env> implements 
Comparable<Procedure<Env>> {
   public int compareTo(Procedure<Env> other) {
     return Long.compare(getProcId(), other.getProcId());
   }
+
+  /**
+   * This function will be called just when procedure is submitted for 
execution.
+   *
+   * @param env The environment passed to the procedure executor
+   */
+  protected void updateMetricsOnSubmit(Env env) {
+    if (env instanceof ConfigNodeProcedureEnv) {
+      ((ConfigNodeProcedureEnv) env)
+          .getConfigManager()
+          .getProcedureManager()
+          .getProcedureMetrics()
+          .updateMetricsOnSubmit(getProcType());
+    }
+  }
+
+  /**
+   * This function will be called just after procedure execution is finished. 
Override this method
+   * to update metrics at the end of the procedure. The default implementation 
adds runtime of a
+   * procedure to a time histogram for successfully completed procedures. 
Increments failed counter
+   * for failed procedures.
+   *
+   * @param env The environment passed to the procedure executor
+   * @param runtime Runtime of the procedure in milliseconds
+   * @param success true if procedure is completed successfully
+   */
+  protected void updateMetricsOnFinish(Env env, long runtime, boolean success) 
{
+    if (env instanceof ConfigNodeProcedureEnv) {
+      ((ConfigNodeProcedureEnv) env)
+          .getConfigManager()
+          .getProcedureManager()
+          .getProcedureMetrics()
+          .updateMetricsOnFinish(getProcType(), runtime, success);
+    }
+  }
 }
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/ProcedureExecutor.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/ProcedureExecutor.java
index d9175797a4e..c9a20080473 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/ProcedureExecutor.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/ProcedureExecutor.java
@@ -391,7 +391,9 @@ public class ProcedureExecutor<Env> {
       rootProcStack.release();
 
       if (proc.isSuccess()) {
-        LOG.info("{} finished in {}ms successfully.", proc, 
proc.elapsedTime());
+        // update metrics on finishing the procedure
+        proc.updateMetricsOnFinish(getEnvironment(), proc.elapsedTime(), true);
+        LOG.debug("{} finished in {}ms successfully.", proc, 
proc.elapsedTime());
         if (proc.getProcId() == rootProcId) {
           rootProcedureCleanup(proc);
         } else {
@@ -509,6 +511,7 @@ public class ProcedureExecutor<Env> {
    */
   private void submitChildrenProcedures(Procedure<Env>[] subprocs) {
     for (Procedure<Env> subproc : subprocs) {
+      subproc.updateMetricsOnSubmit(getEnvironment());
       procedures.put(subproc.getProcId(), subproc);
       scheduler.addFront(subproc);
     }
@@ -662,6 +665,10 @@ public class ProcedureExecutor<Env> {
       if (!procedure.isSuccess()) {
         procedure.setState(ProcedureState.ROLLEDBACK);
       }
+
+      // update metrics on finishing the procedure (fail)
+      procedure.updateMetricsOnFinish(getEnvironment(), 
procedure.elapsedTime(), false);
+
       if (procedure.hasParent()) {
         store.delete(procedure.getProcId());
         procedures.remove(procedure.getProcId());
@@ -709,6 +716,8 @@ public class ProcedureExecutor<Env> {
    */
   private long pushProcedure(Procedure<Env> procedure) {
     final long currentProcId = procedure.getProcId();
+    // Update metrics on start of a procedure
+    procedure.updateMetricsOnSubmit(getEnvironment());
     RootProcedureStack<Env> stack = new RootProcedureStack<>();
     rollbackStack.put(currentProcId, stack);
     procedures.put(currentProcId, procedure);
@@ -797,9 +806,9 @@ public class ProcedureExecutor<Env> {
   }
 
   private final class WorkerMonitor extends InternalProcedure<Env> {
-    private static final int DEFAULT_WORKER_MONITOR_INTERVAL = 5000; // 5sec
+    private static final int DEFAULT_WORKER_MONITOR_INTERVAL = 30000; // 30sec
 
-    private static final int DEFAULT_WORKER_STUCK_THRESHOLD = 10000; // 10sec
+    private static final int DEFAULT_WORKER_STUCK_THRESHOLD = 60000; // 60sec
 
     private static final float DEFAULT_WORKER_ADD_STUCK_PERCENTAGE = 0.5f; // 
50% stuck
 
@@ -854,6 +863,10 @@ public class ProcedureExecutor<Env> {
     return workerThreads.size();
   }
 
+  public long getActiveWorkerThreadCount() {
+    return workerThreads.stream().filter(worker -> 
worker.activeProcedure.get() != null).count();
+  }
+
   public boolean isRunning() {
     return running.get();
   }
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/ProcedureMetrics.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/ProcedureMetrics.java
new file mode 100644
index 00000000000..25533667431
--- /dev/null
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/ProcedureMetrics.java
@@ -0,0 +1,184 @@
+/*
+ * 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.confignode.procedure;
+
+import org.apache.iotdb.commons.service.metric.enums.Metric;
+import org.apache.iotdb.commons.service.metric.enums.Tag;
+import org.apache.iotdb.confignode.manager.ProcedureManager;
+import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
+import org.apache.iotdb.confignode.procedure.scheduler.ProcedureScheduler;
+import org.apache.iotdb.metrics.AbstractMetricService;
+import org.apache.iotdb.metrics.impl.DoNothingMetricManager;
+import org.apache.iotdb.metrics.metricsets.IMetricSet;
+import org.apache.iotdb.metrics.type.Counter;
+import org.apache.iotdb.metrics.type.Timer;
+import org.apache.iotdb.metrics.utils.MetricLevel;
+import org.apache.iotdb.metrics.utils.MetricType;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.ConcurrentHashMap;
+
+public class ProcedureMetrics implements IMetricSet {
+  private static final Logger LOGGER = 
LoggerFactory.getLogger(ProcedureMetrics.class);
+
+  private final ProcedureManager procedureManager;
+
+  private final ProcedureExecutor<ConfigNodeProcedureEnv> executor;
+
+  private final ProcedureScheduler scheduler;
+
+  private final Map<String, ProcedureMetricItems> metricItemsMap = new 
ConcurrentHashMap<>();
+
+  private AbstractMetricService metricService;
+
+  public ProcedureMetrics(ProcedureManager procedureManager) {
+    this.procedureManager = procedureManager;
+    this.executor = procedureManager.getExecutor();
+    this.scheduler = procedureManager.getScheduler();
+  }
+
+  @Override
+  public void bindTo(AbstractMetricService metricService) {
+    this.metricService = metricService;
+    bindThreadMetrics(metricService);
+  }
+
+  @Override
+  public void unbindFrom(AbstractMetricService metricService) {
+    unbindThreadMetrics(metricService);
+
+    for (ProcedureMetricItems items : metricItemsMap.values()) {
+      items.unbindMetricItems(metricService);
+    }
+    metricItemsMap.clear();
+  }
+
+  private void bindThreadMetrics(AbstractMetricService metricService) {
+    metricService.createAutoGauge(
+        Metric.PROCEDURE_WORKER_THREAD_COUNT.toString(),
+        MetricLevel.CORE,
+        executor,
+        ProcedureExecutor::getWorkerThreadCount);
+    metricService.createAutoGauge(
+        Metric.PROCEDURE_ACTIVE_WORKER_THREAD_COUNT.toString(),
+        MetricLevel.CORE,
+        executor,
+        ProcedureExecutor::getActiveWorkerThreadCount);
+    metricService.createAutoGauge(
+        Metric.PROCEDURE_QUEUE_LENGTH.toString(),
+        MetricLevel.CORE,
+        scheduler,
+        ProcedureScheduler::size);
+  }
+
+  private void unbindThreadMetrics(AbstractMetricService metricService) {
+    metricService.remove(MetricType.AUTO_GAUGE, 
Metric.PROCEDURE_WORKER_THREAD_COUNT.toString());
+    metricService.remove(
+        MetricType.AUTO_GAUGE, 
Metric.PROCEDURE_ACTIVE_WORKER_THREAD_COUNT.toString());
+    metricService.remove(MetricType.AUTO_GAUGE, 
Metric.PROCEDURE_QUEUE_LENGTH.toString());
+  }
+
+  public void updateMetricsOnSubmit(String procType) {
+    Optional.ofNullable(metricService)
+        .ifPresent(
+            service -> {
+              metricItemsMap
+                  .computeIfAbsent(procType, k -> new ProcedureMetricItems(k, 
service))
+                  .updateMetricsOnSubmit();
+            });
+  }
+
+  public void updateMetricsOnFinish(String procType, long runtime, boolean 
success) {
+    Optional.ofNullable(metricService)
+        .ifPresent(
+            service -> {
+              metricItemsMap
+                  .computeIfAbsent(procType, k -> new ProcedureMetricItems(k, 
service))
+                  .updateMetricsOnFinish(runtime, success);
+            });
+  }
+
+  class ProcedureMetricItems {
+    private final String procType;
+    private Counter submittedCounter = 
DoNothingMetricManager.DO_NOTHING_COUNTER;
+    private Counter failedCounter = DoNothingMetricManager.DO_NOTHING_COUNTER;
+    private Timer executionTimer = DoNothingMetricManager.DO_NOTHING_TIMER;
+
+    public ProcedureMetricItems(String procType, AbstractMetricService 
metricService) {
+      this.procType = procType;
+      bindMetricItems(metricService);
+    }
+
+    public void bindMetricItems(AbstractMetricService metricService) {
+      this.submittedCounter =
+          metricService.getOrCreateCounter(
+              Metric.PROCEDURE_SUBMITTED_COUNT.toString(),
+              MetricLevel.CORE,
+              Tag.TYPE.toString(),
+              procType);
+      this.failedCounter =
+          metricService.getOrCreateCounter(
+              Metric.PROCEDURE_FAILED_COUNT.toString(),
+              MetricLevel.CORE,
+              Tag.TYPE.toString(),
+              procType);
+      this.executionTimer =
+          metricService.getOrCreateTimer(
+              Metric.PROCEDURE_EXECUTION_TIME.toString(),
+              MetricLevel.CORE,
+              Tag.TYPE.toString(),
+              procType);
+    }
+
+    public void unbindMetricItems(AbstractMetricService metricService) {
+      metricService.remove(
+          MetricType.COUNTER,
+          Metric.PROCEDURE_SUBMITTED_COUNT.toString(),
+          Tag.TYPE.toString(),
+          procType);
+      metricService.remove(
+          MetricType.COUNTER,
+          Metric.PROCEDURE_FAILED_COUNT.toString(),
+          Tag.TYPE.toString(),
+          procType);
+      metricService.remove(
+          MetricType.TIMER,
+          Metric.PROCEDURE_EXECUTION_TIME.toString(),
+          Tag.TYPE.toString(),
+          procType);
+    }
+
+    public void updateMetricsOnSubmit() {
+      submittedCounter.inc();
+    }
+
+    public void updateMetricsOnFinish(long runtime, boolean success) {
+      if (success) {
+        executionTimer.updateMillis(runtime);
+      } else {
+        failedCounter.inc();
+      }
+    }
+  }
+}
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/statemachine/StateMachineProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/statemachine/StateMachineProcedure.java
index 6debdda88b7..05b3915410e 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/statemachine/StateMachineProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/statemachine/StateMachineProcedure.java
@@ -173,7 +173,7 @@ public abstract class StateMachineProcedure<Env, TState> 
extends Procedure<Env>
         setNextState(getStateId(state));
       }
 
-      LOG.info("{} {}; cycles={}", state, this, cycles);
+      LOG.debug("{} {}; cycles={}", state, this, cycles);
       // Keep running count of cycles
       if (getStateId(state) != this.previousState) {
         this.previousState = getStateId(state);
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/TestProcedureExecutor.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/TestProcedureExecutor.java
index 046e4f9ffa9..dce7f2ba5dc 100644
--- 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/TestProcedureExecutor.java
+++ 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/TestProcedureExecutor.java
@@ -99,7 +99,7 @@ public class TestProcedureExecutor extends TestProcedureBase {
   private int waitThreadCount(final int expectedThreads) {
     long startTime = System.currentTimeMillis();
     while (procExecutor.isRunning()
-        && TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis() - 
startTime) <= 30) {
+        && TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis() - 
startTime) <= 180) {
       if (procExecutor.getWorkerThreadCount() == expectedThreads) {
         break;
       }
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/entity/StuckProcedure.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/entity/StuckProcedure.java
index 1db02597c41..0238b6a33bc 100644
--- 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/entity/StuckProcedure.java
+++ 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/entity/StuckProcedure.java
@@ -38,11 +38,11 @@ public class StuckProcedure extends Procedure<TestProcEnv> {
   @Override
   protected Procedure[] execute(final TestProcEnv env) {
     try {
-      if (!latch.tryAcquire(1, 30, TimeUnit.SECONDS)) {
+      if (!latch.tryAcquire(1, 180, TimeUnit.SECONDS)) {
         throw new Exception("waited too long");
       }
 
-      if (!latch.tryAcquire(1, 30, TimeUnit.SECONDS)) {
+      if (!latch.tryAcquire(1, 180, TimeUnit.SECONDS)) {
         throw new Exception("waited too long");
       }
     } catch (Exception e) {
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 61aff528e38..e948c49ce79 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
@@ -38,6 +38,12 @@ public enum Metric {
   TIME_SLOT_NUM_IN_DATABASE("time_slot_num_in_database"),
   REGION_GROUP_NUM_IN_DATABASE("region_group_num_in_database"),
   REPLICATION_FACTOR("replication_factor"),
+  PROCEDURE_WORKER_THREAD_COUNT("procedure_worker_thread_count"),
+  PROCEDURE_ACTIVE_WORKER_THREAD_COUNT("procedure_active_worker_thread_count"),
+  PROCEDURE_QUEUE_LENGTH("procedure_queue_length"),
+  PROCEDURE_SUBMITTED_COUNT("procedure_submitted_count"),
+  PROCEDURE_FAILED_COUNT("procedure_failed_count"),
+  PROCEDURE_EXECUTION_TIME("procedure_execution_time"),
   // protocol related
   ENTRY("entry"),
   SESSION_IDLE_TIME("session_idle_time"),

Reply via email to