From 149a7203bf87905c9d49068a5db9bbb34c172bd1 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Mon, 22 Jun 2026 16:46:48 +0800 Subject: [PATCH] Allow idle Pipe worker threads to time out --- .../concurrent/IoTDBThreadPoolFactory.java | 17 +++++++++++++++++ .../task/execution/PipeSubtaskExecutor.java | 9 ++++++++- .../commons/IoTDBThreadPoolFactoryTest.java | 16 ++++++++++++++++ 3 files changed, 41 insertions(+), 1 deletion(-) diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/IoTDBThreadPoolFactory.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/IoTDBThreadPoolFactory.java index 73d6ff22430d2..14dc807cd2f0f 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/IoTDBThreadPoolFactory.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/IoTDBThreadPoolFactory.java @@ -127,6 +127,23 @@ public static ExecutorService newFixedThreadPool( poolName); } + public static ExecutorService newFixedThreadPoolWithIdleThreadTimeout( + int nThreads, long keepAliveTime, TimeUnit unit, String poolName) { + logger.info(NEW_FIXED_THREAD_POOL_LOGGER_FORMAT, poolName, nThreads); + + WrappedThreadPoolExecutor executor = + new WrappedThreadPoolExecutor( + nThreads, + nThreads, + keepAliveTime, + unit, + new LinkedBlockingQueue<>(), + new IoTThreadFactory(poolName), + poolName); + executor.allowCoreThreadTimeOut(true); + return executor; + } + /** * see {@link Executors#newSingleThreadExecutor(java.util.concurrent.ThreadFactory)}. * diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/execution/PipeSubtaskExecutor.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/execution/PipeSubtaskExecutor.java index c28cfb2569245..7f451729c94bd 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/execution/PipeSubtaskExecutor.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/execution/PipeSubtaskExecutor.java @@ -38,11 +38,14 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; public abstract class PipeSubtaskExecutor { private static final Logger LOGGER = LoggerFactory.getLogger(PipeSubtaskExecutor.class); + private static final long WORKER_THREAD_KEEP_ALIVE_TIME_IN_SECONDS = 60L; + private static final ExecutorService globalSubtaskCallbackListeningExecutor = IoTDBThreadPoolFactory.newSingleThreadExecutor( ThreadName.PIPE_SUBTASK_CALLBACK_EXECUTOR_POOL.getName()); @@ -78,7 +81,11 @@ protected PipeSubtaskExecutor( : ThreadName.PIPE_SUBTASK_CALLBACK_EXECUTOR_POOL.getName(); underlyingThreadPool = (WrappedThreadPoolExecutor) - IoTDBThreadPoolFactory.newFixedThreadPool(corePoolSize, workingThreadName); + IoTDBThreadPoolFactory.newFixedThreadPoolWithIdleThreadTimeout( + corePoolSize, + WORKER_THREAD_KEEP_ALIVE_TIME_IN_SECONDS, + TimeUnit.SECONDS, + workingThreadName); if (disableLogInThreadPool) { underlyingThreadPool.disableErrorLog(); } diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/IoTDBThreadPoolFactoryTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/IoTDBThreadPoolFactoryTest.java index 1db18789fbf69..b735468bed0a3 100644 --- a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/IoTDBThreadPoolFactoryTest.java +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/IoTDBThreadPoolFactoryTest.java @@ -20,6 +20,7 @@ import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.WrappedRunnable; +import org.apache.iotdb.commons.concurrent.threadpool.WrappedThreadPoolExecutor; import org.apache.thrift.server.TThreadPoolServer; import org.apache.thrift.server.TThreadPoolServer.Args; @@ -110,6 +111,21 @@ public void testNewCachedThreadPool() { } } + @Test + public void testNewFixedThreadPoolWithIdleThreadTimeout() throws Exception { + ExecutorService exec = + IoTDBThreadPoolFactory.newFixedThreadPoolWithIdleThreadTimeout( + 1, 1, TimeUnit.MILLISECONDS, POOL_NAME); + + exec.submit(() -> {}).get(); + assertEquals(1, ((WrappedThreadPoolExecutor) exec).getLargestPoolSize()); + + Thread.sleep(100); + assertEquals(0, ((WrappedThreadPoolExecutor) exec).getPoolSize()); + + exec.shutdown(); + } + @Test public void testNewSingleThreadScheduledExecutor() throws InterruptedException { String reason = "(can be ignored in Tests) NewSingleThreadScheduledExecutor";