This is an automated email from the ASF dual-hosted git repository. CRZbulabula pushed a commit to branch yongzao/fix-v2-1039-procedure-metrics-npe in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit b7a2453a73b9a8d23c3a28ac9d57907eb13bffdd Author: Yongzao <[email protected]> AuthorDate: Wed Jul 29 09:51:22 2026 +0800 Fix procedure metrics before executor initialization --- .../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();
