This is an automated email from the ASF dual-hosted git repository.
CRZbulabula 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 7e0812a9881 Fix procedure metrics before executor initialization
(#18338)
7e0812a9881 is described below
commit 7e0812a9881b6f19d85953a570a45ec7c00933eb
Author: Yongzao <[email protected]>
AuthorDate: Wed Jul 29 10:59:01 2026 +0800
Fix procedure metrics before executor initialization (#18338)
---
.../iotdb/confignode/procedure/ProcedureExecutor.java | 12 +++++++++---
.../iotdb/confignode/procedure/TestProcedureExecutor.java | 15 +++++++++++++++
2 files changed, 24 insertions(+), 3 deletions(-)
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 ea2cea8dfad..b5f9a021426 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
@@ -68,7 +68,9 @@ public class ProcedureExecutor<Env> {
private final ThreadGroup threadGroup =
new ThreadGroup(ThreadName.CONFIG_NODE_PROCEDURE_WORKER.getName());
- private CopyOnWriteArrayList<WorkerThread> workerThreads;
+ // Metrics may be scraped before init() and concurrently with initialization
during a ConfigNode
+ // leader transition.
+ private volatile CopyOnWriteArrayList<WorkerThread> workerThreads;
private TimeoutExecutorThread<Env> timeoutExecutor;
@@ -1026,11 +1028,15 @@ public class ProcedureExecutor<Env> {
}
public int getWorkerThreadCount() {
- return workerThreads.size();
+ final CopyOnWriteArrayList<WorkerThread> workers = workerThreads;
+ return workers == null ? 0 : workers.size();
}
public long getActiveWorkerThreadCount() {
- return workerThreads.stream().filter(worker ->
worker.activeProcedure.get() != null).count();
+ final CopyOnWriteArrayList<WorkerThread> workers = workerThreads;
+ return workers == null
+ ? 0
+ : workers.stream().filter(worker -> worker.activeProcedure.get() !=
null).count();
}
public boolean isRunning() {
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 692317ffa2b..04f051e06ee 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
@@ -72,6 +72,21 @@ public class TestProcedureExecutor extends TestProcedureBase
{
Assert.assertTrue(procExecutor.isFinished(procId));
}
+ @Test
+ public void testWorkerThreadMetricsBeforeInitialization() {
+ TestProcEnv localEnv = new TestProcEnv();
+ ProcedureExecutor<TestProcEnv> localExecutor =
+ new ProcedureExecutor<>(localEnv, new NoopProcedureStore());
+ localEnv.setScheduler(localExecutor.getScheduler());
+
+ Assert.assertEquals(0, localExecutor.getWorkerThreadCount());
+ Assert.assertEquals(0, localExecutor.getActiveWorkerThreadCount());
+
+ localExecutor.init(2);
+ Assert.assertEquals(2, localExecutor.getWorkerThreadCount());
+ Assert.assertEquals(0, localExecutor.getActiveWorkerThreadCount());
+ }
+
@Test
public void testProcedureFailedDuringSubmissionIsRolledBack() throws
InterruptedException {
TestProcEnv localEnv = new TestProcEnv();