This is an automated email from the ASF dual-hosted git repository. hxd pushed a commit to branch cluster- in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 248ce36e1549f745b54ba56a83db8e55e9335ea3 Author: xiangdong huang <[email protected]> AuthorDate: Wed Aug 11 08:40:19 2021 +0800 modify thread pool --- .../db/concurrent/IoTDBDaemonThreadFactory.java | 36 ++++++++++++++++++++++ .../db/concurrent/IoTDBThreadPoolFactory.java | 12 ++++++++ .../apache/iotdb/db/cq/ContinuousQueryService.java | 4 +-- .../org/apache/iotdb/db/engine/StorageEngine.java | 6 ++-- .../iotdb/db/engine/merge/manage/MergeManager.java | 6 ++-- .../engine/storagegroup/StorageGroupProcessor.java | 5 +-- .../org/apache/iotdb/db/metadata/MManager.java | 22 ++++++++++--- .../apache/iotdb/db/service/MetricsService.java | 4 +-- .../org/apache/iotdb/db/service/TSServiceImpl.java | 4 +-- .../org/apache/iotdb/db/service/UpgradeSevice.java | 6 ++-- .../writelog/manager/MultiFileLogNodeManager.java | 5 +-- .../db/writelog/node/ExclusiveWriteLogNode.java | 6 ++-- 12 files changed, 87 insertions(+), 29 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/concurrent/IoTDBDaemonThreadFactory.java b/server/src/main/java/org/apache/iotdb/db/concurrent/IoTDBDaemonThreadFactory.java new file mode 100644 index 0000000..4c14e73 --- /dev/null +++ b/server/src/main/java/org/apache/iotdb/db/concurrent/IoTDBDaemonThreadFactory.java @@ -0,0 +1,36 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iotdb.db.concurrent; + +public class IoTDBDaemonThreadFactory extends IoTThreadFactory { + public IoTDBDaemonThreadFactory(String poolName) { + super(poolName); + } + + public IoTDBDaemonThreadFactory(String poolName, Thread.UncaughtExceptionHandler handler) { + super(poolName, handler); + } + + @Override + public Thread newThread(Runnable r) { + Thread t = super.newThread(r); + t.setDaemon(true); + return t; + } +} diff --git a/server/src/main/java/org/apache/iotdb/db/concurrent/IoTDBThreadPoolFactory.java b/server/src/main/java/org/apache/iotdb/db/concurrent/IoTDBThreadPoolFactory.java index 8e4eed9..77799eb 100644 --- a/server/src/main/java/org/apache/iotdb/db/concurrent/IoTDBThreadPoolFactory.java +++ b/server/src/main/java/org/apache/iotdb/db/concurrent/IoTDBThreadPoolFactory.java @@ -127,6 +127,18 @@ public class IoTDBThreadPoolFactory { poolName); } + public static ExecutorService newCachedThreadPoolWithDaemon(String poolName) { + logger.info("new cached thread pool: {}", poolName); + return new WrappedThreadPoolExecutor( + 0, + Integer.MAX_VALUE, + 60L, + TimeUnit.SECONDS, + new SynchronousQueue<>(), + new IoTDBDaemonThreadFactory(poolName), + poolName); + } + /** * see {@link Executors#newSingleThreadExecutor(java.util.concurrent.ThreadFactory)}. * diff --git a/server/src/main/java/org/apache/iotdb/db/cq/ContinuousQueryService.java b/server/src/main/java/org/apache/iotdb/db/cq/ContinuousQueryService.java index eb2f908..fb7a1b2 100644 --- a/server/src/main/java/org/apache/iotdb/db/cq/ContinuousQueryService.java +++ b/server/src/main/java/org/apache/iotdb/db/cq/ContinuousQueryService.java @@ -19,6 +19,7 @@ package org.apache.iotdb.db.cq; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.exception.ContinuousQueryException; import org.apache.iotdb.db.exception.ShutdownException; @@ -36,7 +37,6 @@ import org.slf4j.LoggerFactory; import java.util.ArrayList; import java.util.List; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.ReentrantLock; @@ -94,7 +94,7 @@ public class ContinuousQueryService implements IService { nextExecutionTimestamps.put(plan.getContinuousQueryName(), nextExecutionTimestamp); } - checkThread = Executors.newSingleThreadScheduledExecutor(); + checkThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("CQ-Check"); checkThread.scheduleAtFixedRate( this::checkAndSubmitTasks, 0, diff --git a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java index ef955f5..b38eab5 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java @@ -89,7 +89,6 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @@ -272,13 +271,14 @@ public class StorageEngine implements IService { @Override public void start() { - ttlCheckThread = Executors.newSingleThreadScheduledExecutor(); + ttlCheckThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("TTL"); ttlCheckThread.scheduleAtFixedRate( this::checkTTL, TTL_CHECK_INTERVAL, TTL_CHECK_INTERVAL, TimeUnit.MILLISECONDS); logger.info("start ttl check thread successfully."); if (config.isEnableTimedFlushUnseqMemtable()) { - memtableTimedFlushCheckThread = Executors.newSingleThreadScheduledExecutor(); + memtableTimedFlushCheckThread = + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("TimedFlushMemtable"); memtableTimedFlushCheckThread.scheduleAtFixedRate( this::timedFlushMemTable, config.getUnseqMemtableFlushCheckInterval(), diff --git a/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeManager.java b/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeManager.java index f9112e1..aa478bc 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeManager.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeManager.java @@ -19,6 +19,7 @@ package org.apache.iotdb.db.engine.merge.manage; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.concurrent.ThreadName; import org.apache.iotdb.db.conf.IoTDBConstant; import org.apache.iotdb.db.conf.IoTDBDescriptor; @@ -45,7 +46,6 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentSkipListSet; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadPoolExecutor; @@ -146,13 +146,13 @@ public class MergeManager implements IService, MergeManagerMBean { long mergeInterval = IoTDBDescriptor.getInstance().getConfig().getMergeIntervalSec(); if (mergeInterval > 0) { timedMergeThreadPool = - Executors.newSingleThreadScheduledExecutor(r -> new Thread(r, "TimedMergeThread")); + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("TimedMergeThread"); timedMergeThreadPool.scheduleAtFixedRate( this::mergeAll, mergeInterval, mergeInterval, TimeUnit.SECONDS); } taskCleanerThreadPool = - Executors.newSingleThreadScheduledExecutor(r -> new Thread(r, "MergeTaskCleaner")); + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("MergeTaskCleaner"); taskCleanerThreadPool.scheduleAtFixedRate(this::cleanFinishedTask, 30, 30, TimeUnit.MINUTES); logger.info("MergeManager started"); } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java index 77d68b1..6c52463 100755 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java @@ -18,6 +18,7 @@ */ package org.apache.iotdb.db.engine.storagegroup; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBConstant; import org.apache.iotdb.db.conf.IoTDBDescriptor; @@ -104,7 +105,6 @@ import java.util.Map; import java.util.Map.Entry; import java.util.Set; import java.util.TreeMap; -import java.util.concurrent.Executors; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @@ -406,7 +406,8 @@ public class StorageGroupProcessor { .getCompactionStrategy() .getTsFileManagement(logicalStorageGroupName, storageGroupSysDir.getAbsolutePath()); - ScheduledExecutorService executorService = Executors.newSingleThreadScheduledExecutor(); + ScheduledExecutorService executorService = + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("WAL-Trim"); executorService.scheduleWithFixedDelay( this::trimTask, config.getWalPoolTrimIntervalInMS(), diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java index 6e585cf..8e374ec 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java @@ -18,15 +18,29 @@ */ package org.apache.iotdb.db.metadata; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.engine.StorageEngine; import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory; import org.apache.iotdb.db.engine.trigger.executor.TriggerEngine; -import org.apache.iotdb.db.exception.metadata.*; +import org.apache.iotdb.db.exception.metadata.AliasAlreadyExistException; +import org.apache.iotdb.db.exception.metadata.AlignedTimeseriesException; +import org.apache.iotdb.db.exception.metadata.DataTypeMismatchException; +import org.apache.iotdb.db.exception.metadata.DeleteFailedException; +import org.apache.iotdb.db.exception.metadata.IllegalPathException; +import org.apache.iotdb.db.exception.metadata.MetadataException; +import org.apache.iotdb.db.exception.metadata.PathAlreadyExistException; +import org.apache.iotdb.db.exception.metadata.PathNotExistException; +import org.apache.iotdb.db.exception.metadata.StorageGroupAlreadySetException; +import org.apache.iotdb.db.exception.metadata.StorageGroupNotSetException; import org.apache.iotdb.db.metadata.logfile.MLogReader; import org.apache.iotdb.db.metadata.logfile.MLogWriter; -import org.apache.iotdb.db.metadata.mnode.*; +import org.apache.iotdb.db.metadata.mnode.IEntityMNode; +import org.apache.iotdb.db.metadata.mnode.IMNode; +import org.apache.iotdb.db.metadata.mnode.IMeasurementMNode; +import org.apache.iotdb.db.metadata.mnode.IStorageGroupMNode; +import org.apache.iotdb.db.metadata.mnode.MeasurementMNode; import org.apache.iotdb.db.metadata.tag.TagManager; import org.apache.iotdb.db.metadata.template.Template; import org.apache.iotdb.db.metadata.template.TemplateManager; @@ -91,7 +105,6 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Set; -import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; @@ -184,8 +197,7 @@ public class MManager { if (config.isEnableMTreeSnapshot()) { timedCreateMTreeSnapshotThread = - Executors.newSingleThreadScheduledExecutor( - r -> new Thread(r, "timedCreateMTreeSnapshotThread")); + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("timedCreateMTreeSnapshot"); timedCreateMTreeSnapshotThread.scheduleAtFixedRate( this::checkMTreeModified, MTREE_SNAPSHOT_THREAD_CHECK_TIME, diff --git a/server/src/main/java/org/apache/iotdb/db/service/MetricsService.java b/server/src/main/java/org/apache/iotdb/db/service/MetricsService.java index cb15eac..6925416 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/MetricsService.java +++ b/server/src/main/java/org/apache/iotdb/db/service/MetricsService.java @@ -18,6 +18,7 @@ */ package org.apache.iotdb.db.service; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.concurrent.ThreadName; import org.apache.iotdb.db.concurrent.WrappedRunnable; import org.apache.iotdb.db.conf.IoTDBConfig; @@ -37,7 +38,6 @@ import java.net.InetSocketAddress; import java.net.Socket; import java.net.SocketAddress; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; public class MetricsService implements MetricsServiceMBean, IService { @@ -90,7 +90,7 @@ public class MetricsService implements MetricsServiceMBean, IService { return; } logger.info("{}: start {}...", IoTDBConstant.GLOBAL_DB_NAME, this.getID().getName()); - executorService = Executors.newSingleThreadExecutor(); + executorService = IoTDBThreadPoolFactory.newSingleThreadExecutor("Deprecated-Metric-Monitor"); int port = getMetricsPort(); MetricsSystem metricsSystem = new MetricsSystem(new ServerArgument(port)); MetricsWebUI metricsWebUI = new MetricsWebUI(metricsSystem.getMetricRegistry()); diff --git a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java index 59ec25f..c2b4009 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java @@ -22,6 +22,7 @@ import org.apache.iotdb.db.auth.AuthException; import org.apache.iotdb.db.auth.AuthorityChecker; import org.apache.iotdb.db.auth.authorizer.BasicAuthorizer; import org.apache.iotdb.db.auth.authorizer.IAuthorizer; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBConstant; import org.apache.iotdb.db.conf.IoTDBDescriptor; @@ -150,7 +151,6 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Set; -import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -199,7 +199,7 @@ public class TSServiceImpl implements TSIService.Iface { executor = new PlanExecutor(); ScheduledExecutorService timedQuerySqlCountThread = - Executors.newSingleThreadScheduledExecutor(r -> new Thread(r, "timedQuerySqlCountThread")); + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("timedQuerySqlCount"); timedQuerySqlCountThread.scheduleAtFixedRate( () -> { if (queryCount.get() != 0) { diff --git a/server/src/main/java/org/apache/iotdb/db/service/UpgradeSevice.java b/server/src/main/java/org/apache/iotdb/db/service/UpgradeSevice.java index faa3c39..bb8558c 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/UpgradeSevice.java +++ b/server/src/main/java/org/apache/iotdb/db/service/UpgradeSevice.java @@ -18,6 +18,7 @@ */ package org.apache.iotdb.db.service; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.engine.StorageEngine; import org.apache.iotdb.db.engine.upgrade.UpgradeLog; @@ -29,7 +30,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.atomic.AtomicInteger; public class UpgradeSevice implements IService { @@ -58,9 +58,7 @@ public class UpgradeSevice implements IService { if (updateThreadNum <= 0) { updateThreadNum = 1; } - upgradeThreadPool = - Executors.newFixedThreadPool( - updateThreadNum, r -> new Thread(r, "UpgradeThread-" + threadCnt.getAndIncrement())); + upgradeThreadPool = IoTDBThreadPoolFactory.newFixedThreadPool(updateThreadNum, "UpgradeThread"); UpgradeLog.createUpgradeLog(); countUpgradeFiles(); if (cntUpgradeFileNum.get() == 0) { diff --git a/server/src/main/java/org/apache/iotdb/db/writelog/manager/MultiFileLogNodeManager.java b/server/src/main/java/org/apache/iotdb/db/writelog/manager/MultiFileLogNodeManager.java index d432c61..236ef2f 100644 --- a/server/src/main/java/org/apache/iotdb/db/writelog/manager/MultiFileLogNodeManager.java +++ b/server/src/main/java/org/apache/iotdb/db/writelog/manager/MultiFileLogNodeManager.java @@ -18,6 +18,7 @@ */ package org.apache.iotdb.db.writelog.manager; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.exception.StartupException; @@ -33,7 +34,6 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.function.Consumer; @@ -122,7 +122,8 @@ public class MultiFileLogNodeManager implements WriteLogNodeManager, IService { return; } if (config.getForceWalPeriodInMs() > 0) { - executorService = Executors.newSingleThreadScheduledExecutor(); + executorService = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("WAL-ForceSync"); + executorService.scheduleWithFixedDelay( this::forceTask, config.getForceWalPeriodInMs(), diff --git a/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java b/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java index c947963..e59aada 100644 --- a/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java +++ b/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java @@ -18,6 +18,7 @@ */ package org.apache.iotdb.db.writelog.node; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.conf.directories.DirectoryManager; @@ -28,7 +29,6 @@ import org.apache.iotdb.db.writelog.io.ILogWriter; import org.apache.iotdb.db.writelog.io.LogWriter; import org.apache.iotdb.db.writelog.io.MultiFileLogReader; -import com.google.common.util.concurrent.ThreadFactoryBuilder; import org.apache.commons.io.FileUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -42,7 +42,6 @@ import java.nio.channels.ClosedChannelException; import java.util.Arrays; import java.util.Comparator; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.locks.ReentrantLock; /** This WriteLogNode is used to manage insert ahead logs of a TsFile. */ @@ -69,8 +68,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode, Comparable<Exclusive private final Object switchBufferCondition = new Object(); private ReentrantLock lock = new ReentrantLock(); private static final ExecutorService FLUSH_BUFFER_THREAD_POOL = - Executors.newCachedThreadPool( - new ThreadFactoryBuilder().setNameFormat("Flush-WAL-Thread-%d").setDaemon(true).build()); + IoTDBThreadPoolFactory.newCachedThreadPoolWithDaemon("Flush-WAL-Thread"); private long fileId = 0; private long lastFlushedId = 0;
