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