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();

Reply via email to