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"),