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

jt2594838 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 a22746a57af Fix thread pool lifecycle leaks (#18703)
a22746a57af is described below

commit a22746a57af734b23a7a2b91567f4f75ae1140b4
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 24 17:09:13 2026 +0800

    Fix thread pool lifecycle leaks (#18703)
---
 .../db/partition/DataPartitionTableGenerator.java  |  8 +++-
 .../impl/DataNodeInternalRPCServiceImpl.java       | 50 ++++++++++++----------
 .../metrics/IoTDBInternalLocalReporter.java        |  1 +
 .../iotdb/db/storageengine/StorageEngine.java      |  4 ++
 .../db/storageengine/dataregion/DataRegion.java    |  7 +++
 .../reporter/iotdb/IoTDBSessionReporter.java       |  1 +
 6 files changed, 47 insertions(+), 24 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/partition/DataPartitionTableGenerator.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/partition/DataPartitionTableGenerator.java
index ed85dda0761..2f84fcbedbf 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/partition/DataPartitionTableGenerator.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/partition/DataPartitionTableGenerator.java
@@ -120,7 +120,13 @@ public class DataPartitionTableGenerator {
     }
 
     status = TaskStatus.IN_PROGRESS;
-    return 
CompletableFuture.runAsync(this::generateDataPartitionTableByMemory);
+    return CompletableFuture.runAsync(this::generateDataPartitionTableByMemory)
+        .whenComplete((ignored, throwable) -> close());
+  }
+
+  /** Close the executor owned by this generator after generation has 
finished. */
+  public void close() {
+    executor.shutdownNow();
   }
 
   private void generateDataPartitionTableByMemory() {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
index e7213b7822c..8f2c1080b3f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
@@ -3789,31 +3789,35 @@ public class DataNodeInternalRPCServiceImpl implements 
IDataNodeRPCService.Iface
             ThreadName.FIND_EARLIEST_TIME_SLOT_PARALLEL_POOL.getName(),
             new ThreadPoolExecutor.CallerRunsPolicy());
 
-    for (DataRegion dataRegion : 
StorageEngine.getInstance().getAllDataRegions()) {
-      CompletableFuture<Void> regionFuture =
-          CompletableFuture.runAsync(
-              () -> {
-                TsFileManager tsFileManager = dataRegion.getTsFileManager();
-                String databaseName = dataRegion.getDatabaseName();
-                if (ignoreDatabase.contains(databaseName)) {
-                  return;
-                }
+    try {
+      for (DataRegion dataRegion : 
StorageEngine.getInstance().getAllDataRegions()) {
+        CompletableFuture<Void> regionFuture =
+            CompletableFuture.runAsync(
+                () -> {
+                  TsFileManager tsFileManager = dataRegion.getTsFileManager();
+                  String databaseName = dataRegion.getDatabaseName();
+                  if (ignoreDatabase.contains(databaseName)) {
+                    return;
+                  }
 
-                Set<Long> timePartitionIds = tsFileManager.getTimePartitions();
-                if (timePartitionIds.isEmpty()) {
-                  return;
-                }
-                final long earliestTimeSlotId = 
Collections.min(timePartitionIds);
-                earliestTimeslots.compute(
-                    databaseName,
-                    (k, v) -> v == null ? earliestTimeSlotId : 
Math.min(earliestTimeSlotId, v));
-              },
-              findEarliestTimeSlotExecutor);
-      futures.add(regionFuture);
-    }
+                  Set<Long> timePartitionIds = 
tsFileManager.getTimePartitions();
+                  if (timePartitionIds.isEmpty()) {
+                    return;
+                  }
+                  final long earliestTimeSlotId = 
Collections.min(timePartitionIds);
+                  earliestTimeslots.compute(
+                      databaseName,
+                      (k, v) -> v == null ? earliestTimeSlotId : 
Math.min(earliestTimeSlotId, v));
+                },
+                findEarliestTimeSlotExecutor);
+        futures.add(regionFuture);
+      }
 
-    // Wait for all tasks to complete
-    CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
+      // Wait for all tasks to complete
+      CompletableFuture.allOf(futures.toArray(new 
CompletableFuture[0])).join();
+    } finally {
+      findEarliestTimeSlotExecutor.shutdown();
+    }
     LOGGER.info(DataNodeMiscMessages.PROCESS_DATA_DIR_COMPLETED);
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java
index f081a7b60fc..1630175738f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java
@@ -153,6 +153,7 @@ public class IoTDBInternalLocalReporter extends 
IoTDBInternalReporter {
       currentServiceFuture.cancel(true);
       currentServiceFuture = null;
     }
+    service.shutdownNow();
     clear();
     LOGGER.info(DataNodeMiscMessages.INTERNAL_REPORTER_STOP);
     return true;
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
index a06ec8c00b5..2e779717624 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
@@ -440,6 +440,7 @@ public class StorageEngine implements IService {
       }
     }
     syncCloseAllProcessor();
+    
dataRegionMap.values().forEach(DataRegion::shutdownUpgradeModFileThreadPool);
     ThreadUtils.stopThreadPool(
         seqMemtableTimedFlushCheckThread, ThreadName.TIMED_FLUSH_SEQ_MEMTABLE);
     ThreadUtils.stopThreadPool(
@@ -462,6 +463,8 @@ public class StorageEngine implements IService {
       forceCloseAllProcessor();
     } catch (TsFileProcessorException e) {
       throw new ShutdownException(e);
+    } finally {
+      
dataRegionMap.values().forEach(DataRegion::shutdownUpgradeModFileThreadPool);
     }
     shutdownTimedService(seqMemtableTimedFlushCheckThread, 
"SeqMemtableTimedFlushCheckThread");
     shutdownTimedService(unseqMemtableTimedFlushCheckThread, 
"UnseqMemtableTimedFlushCheckThread");
@@ -517,6 +520,7 @@ public class StorageEngine implements IService {
   /** This function is just for unit test. */
   @TestOnly
   public synchronized void reset() {
+    
dataRegionMap.values().forEach(DataRegion::shutdownUpgradeModFileThreadPool);
     dataRegionMap.clear();
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
index 45c3e4d560f..e8da377db6b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
@@ -2329,6 +2329,12 @@ public class DataRegion implements IDataRegionForQuery {
     }
   }
 
+  public void shutdownUpgradeModFileThreadPool() {
+    if (upgradeModFileThreadPool != null) {
+      upgradeModFileThreadPool.shutdown();
+    }
+  }
+
   public void deleteDALFolderAndClose() {
     
Optional.ofNullable(DeletionResourceManager.getInstance(dataRegionId.getId()))
         .ifPresent(
@@ -5262,6 +5268,7 @@ public class DataRegion implements IDataRegionForQuery {
     writeLock("markDeleted");
     try {
       deleted = true;
+      shutdownUpgradeModFileThreadPool();
       releaseDirectBufferMemory();
       MetricService.getInstance().removeMetricSet(metrics);
       deletedCondition.signalAll();
diff --git 
a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java
 
b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java
index b4246ad9e85..48a22f168dd 100644
--- 
a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java
+++ 
b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java
@@ -137,6 +137,7 @@ public class IoTDBSessionReporter extends IoTDBReporter {
       currentServiceFuture.cancel(true);
       currentServiceFuture = null;
     }
+    service.shutdownNow();
     if (sessionPool != null) {
       sessionPool.close();
     }

Reply via email to