Skip to content

Commit d3cebd2

Browse files
authored
[To dev/1.3] Fix thread pool lifecycle leaks (#18703) (#18720)
* Fix thread pool lifecycle leaks (#18703) * Adapt thread pool cherry-pick for dev/1.3
1 parent e3f50df commit d3cebd2

4 files changed

Lines changed: 36 additions & 24 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/partition/DataPartitionTableGenerator.java‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -120,7 +120,13 @@ public CompletableFuture<Void> startGeneration() {
120120
}
121121

122122
status = TaskStatus.IN_PROGRESS;
123-
return CompletableFuture.runAsync(this::generateDataPartitionTableByMemory);
123+
return CompletableFuture.runAsync(this::generateDataPartitionTableByMemory)
124+
.whenComplete((ignored, throwable) -> close());
125+
}
126+
127+
/** Close the executor owned by this generator after generation has finished. */
128+
public void close() {
129+
executor.shutdownNow();
124130
}
125131

126132
private void generateDataPartitionTableByMemory() {

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java‎

Lines changed: 27 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -2916,31 +2916,35 @@ private void processDataRegionForEarliestTimeslots(Map<String, Long> earliestTim
29162916
ThreadName.FIND_EARLIEST_TIME_SLOT_PARALLEL_POOL.getName(),
29172917
new ThreadPoolExecutor.CallerRunsPolicy());
29182918

2919-
for (DataRegion dataRegion : StorageEngine.getInstance().getAllDataRegions()) {
2920-
CompletableFuture<Void> regionFuture =
2921-
CompletableFuture.runAsync(
2922-
() -> {
2923-
TsFileManager tsFileManager = dataRegion.getTsFileManager();
2924-
String databaseName = dataRegion.getDatabaseName();
2925-
if (ignoreDatabase.contains(databaseName)) {
2926-
return;
2927-
}
2919+
try {
2920+
for (DataRegion dataRegion : StorageEngine.getInstance().getAllDataRegions()) {
2921+
CompletableFuture<Void> regionFuture =
2922+
CompletableFuture.runAsync(
2923+
() -> {
2924+
TsFileManager tsFileManager = dataRegion.getTsFileManager();
2925+
String databaseName = dataRegion.getDatabaseName();
2926+
if (ignoreDatabase.contains(databaseName)) {
2927+
return;
2928+
}
29282929

2929-
Set<Long> timePartitionIds = tsFileManager.getTimePartitions();
2930-
if (timePartitionIds.isEmpty()) {
2931-
return;
2932-
}
2933-
final long earliestTimeSlotId = Collections.min(timePartitionIds);
2934-
earliestTimeslots.compute(
2935-
databaseName,
2936-
(k, v) -> v == null ? earliestTimeSlotId : Math.min(earliestTimeSlotId, v));
2937-
},
2938-
findEarliestTimeSlotExecutor);
2939-
futures.add(regionFuture);
2940-
}
2930+
Set<Long> timePartitionIds = tsFileManager.getTimePartitions();
2931+
if (timePartitionIds.isEmpty()) {
2932+
return;
2933+
}
2934+
final long earliestTimeSlotId = Collections.min(timePartitionIds);
2935+
earliestTimeslots.compute(
2936+
databaseName,
2937+
(k, v) -> v == null ? earliestTimeSlotId : Math.min(earliestTimeSlotId, v));
2938+
},
2939+
findEarliestTimeSlotExecutor);
2940+
futures.add(regionFuture);
2941+
}
29412942

2942-
// Wait for all tasks to complete
2943-
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
2943+
// Wait for all tasks to complete
2944+
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
2945+
} finally {
2946+
findEarliestTimeSlotExecutor.shutdown();
2947+
}
29442948
LOGGER.info("Process data directory for earliestTimeslots completed successfully");
29452949
}
29462950

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -148,6 +148,7 @@ public boolean stop() {
148148
currentServiceFuture.cancel(true);
149149
currentServiceFuture = null;
150150
}
151+
service.shutdownNow();
151152
clear();
152153
LOGGER.info("IoTDBInternalReporter stop!");
153154
return true;

‎iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -128,6 +128,7 @@ public boolean stop() {
128128
currentServiceFuture.cancel(true);
129129
currentServiceFuture = null;
130130
}
131+
service.shutdownNow();
131132
if (sessionPool != null) {
132133
sessionPool.close();
133134
}

0 commit comments

Comments
 (0)