This is an automated email from the ASF dual-hosted git repository. ericpai pushed a commit to branch enhance/iotdb-3407 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 6685a93cbff4e0f70697711f53e3099b4a9f93e3 Author: ericpai <[email protected]> AuthorDate: Tue Jun 7 14:07:53 2022 +0800 [IOTDB-3407] Checkstyle: force to use safe thread schedule interface --- checkstyle.xml | 30 +++--- .../org/apache/iotdb/cluster/ClusterIoTDB.java | 7 +- .../iotdb/cluster/log/manage/RaftLogManager.java | 5 +- .../serializable/SyncLogDequeSerializer.java | 4 +- .../cluster/server/PullSnapshotHintService.java | 4 +- .../cluster/utils/nodetool/ClusterMonitor.java | 28 +++--- .../iotdb/confignode/manager/PartitionManager.java | 4 +- .../java/org/apache/iotdb/flink/IoTDBSink.java | 1 + .../threadpool/ScheduledExecutorUtil.java | 106 +++++++++++++++++++++ .../WrappedScheduledExecutorService.java | 2 + .../WrappedSingleThreadScheduledExecutor.java | 2 + .../iotdb/commons/IoTDBThreadPoolFactoryTest.java | 7 +- pom.xml | 4 +- .../org/apache/iotdb/db/engine/StorageEngine.java | 27 ++++-- .../apache/iotdb/db/engine/StorageEngineV2.java | 15 ++- .../engine/compaction/CompactionTaskManager.java | 12 +-- .../iotdb/db/engine/cq/ContinuousQueryService.java | 4 +- .../iotdb/db/engine/storagegroup/DataRegion.java | 24 +++-- .../iotdb/db/localconfignode/LocalConfigNode.java | 5 +- .../fragment/FragmentInstanceManager.java | 14 +-- .../scheduler/FixedRateFragInsStateTracker.java | 9 +- .../db/query/control/SessionTimeoutManager.java | 4 +- .../db/service/basic/QueryFrequencyRecorder.java | 4 +- .../db/sync/sender/service/TransportHandler.java | 15 +-- .../java/org/apache/iotdb/db/wal/WALManager.java | 14 +-- 25 files changed, 249 insertions(+), 102 deletions(-) diff --git a/checkstyle.xml b/checkstyle.xml index 4d4eb175a9..5af49cd7e4 100644 --- a/checkstyle.xml +++ b/checkstyle.xml @@ -40,7 +40,25 @@ <module name="FileTabCharacter"> <property name="eachLine" value="true"/> </module> + <module name="LineLength"> + <property name="max" value="100"/> + <property name="ignorePattern" value="^package.*|^import.*|a href|href|http://|https://|ftp://"/> + </module> + <module name="SuppressWarningsFilter" /> <module name="TreeWalker"> + <module name="SuppressWarningsHolder" /> + <!--ERROR severity rules, each Java file should obey this --> + <module name="AvoidStarImport"> + <property name="severity" value="error"/> + </module> + <module name="RegexpSinglelineJava"> + <property name="id" value="unsafeThreadSchedule"/> + <property name="format" value="schedule(AtFixedRate|WithFixedDelay)\("/> + <property name="severity" value="error"/> + <property name="ignoreComments" value="true"/> + <property name="message" value="Using ScheduledExecutorService::schedule(AtFixedRate|WithFixedDelay) directly is unsafe, please use ScheduledExecutorUtil::safelySchedule(AtFixedRate|WithFixedDelay) instead."/> + </module> + <!--WARNING severity rules, which may be promoted to error level--> <module name="OuterTypeFilename"/> <module name="IllegalTokenText"> <property name="tokens" value="STRING_LITERAL, CHAR_LITERAL"/> @@ -52,13 +70,6 @@ <property name="allowByTailComment" value="true"/> <property name="allowNonPrintableEscapes" value="true"/> </module> - <module name="LineLength"> - <property name="max" value="100"/> - <property name="ignorePattern" value="^package.*|^import.*|a href|href|http://|https://|ftp://"/> - </module> - <module name="AvoidStarImport"> - <property name="severity" value="error"/> - </module> <module name="OneTopLevelClass"/> <module name="NoLineWrap"/> <module name="EmptyBlock"> @@ -212,13 +223,10 @@ <property name="target" value="CLASS_DEF, INTERFACE_DEF, ENUM_DEF, METHOD_DEF, CTOR_DEF, VARIABLE_DEF"/> </module> <module name="JavadocMethod"> - <property name="scope" value="public"/> <property name="allowMissingParamTags" value="true"/> - <property name="allowMissingThrowsTags" value="true"/> + <property name="validateThrows" value="true"/> <property name="allowMissingReturnTag" value="true"/> - <property name="minLineCount" value="2"/> <property name="allowedAnnotations" value="Override, Test"/> - <property name="allowThrowsTagsForSubclasses" value="true"/> </module> <module name="MethodName"> <property name="format" value="^[a-z][a-z0-9][a-zA-Z0-9_]*$"/> diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java index 2910a3fd47..e4ddfd87b2 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java @@ -57,6 +57,7 @@ import org.apache.iotdb.cluster.server.service.MetaSyncService; import org.apache.iotdb.cluster.utils.ClusterUtils; import org.apache.iotdb.cluster.utils.nodetool.ClusterMonitor; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.exception.ConfigurationException; import org.apache.iotdb.commons.exception.StartupException; @@ -189,14 +190,16 @@ public class ClusterIoTDB implements ClusterIoTDBMBean { private void initTasks() { reportThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("NodeReportThread"); - reportThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + reportThread, this::generateNodeReport, ClusterConstant.REPORT_INTERVAL_SEC, ClusterConstant.REPORT_INTERVAL_SEC, TimeUnit.SECONDS); hardLinkCleanerThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("HardLinkCleaner"); - hardLinkCleanerThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + hardLinkCleanerThread, new HardLinkCleaner(), ClusterConstant.CLEAN_HARDLINK_INTERVAL_SEC, ClusterConstant.CLEAN_HARDLINK_INTERVAL_SEC, diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/RaftLogManager.java b/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/RaftLogManager.java index 3a675afa4f..d23bfdded0 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/RaftLogManager.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/RaftLogManager.java @@ -33,6 +33,7 @@ import org.apache.iotdb.cluster.log.StableEntryManager; import org.apache.iotdb.cluster.log.manage.serializable.LogManagerMeta; import org.apache.iotdb.cluster.server.monitor.Timer.Statistic; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.utils.TestOnly; import org.apache.iotdb.tsfile.utils.RamUsageEstimator; @@ -152,8 +153,10 @@ public abstract class RaftLogManager { ClusterDescriptor.getInstance().getConfig().getLogDeleteCheckIntervalSecond(); if (logDeleteCheckIntervalSecond > 0) { + this.deleteLogFuture = - deleteLogExecutorService.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + deleteLogExecutorService, this::checkDeleteLog, logDeleteCheckIntervalSecond, logDeleteCheckIntervalSecond, diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/serializable/SyncLogDequeSerializer.java b/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/serializable/SyncLogDequeSerializer.java index beacdf5c23..655086a489 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/serializable/SyncLogDequeSerializer.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/serializable/SyncLogDequeSerializer.java @@ -24,6 +24,7 @@ import org.apache.iotdb.cluster.log.HardState; import org.apache.iotdb.cluster.log.Log; import org.apache.iotdb.cluster.log.LogParser; import org.apache.iotdb.cluster.log.StableEntryManager; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.file.SystemFileFactory; import org.apache.iotdb.commons.utils.TestOnly; import org.apache.iotdb.db.conf.IoTDBDescriptor; @@ -169,7 +170,8 @@ public class SyncLogDequeSerializer implements StableEntryManager { .build()); this.persistLogDeleteLogFuture = - persistLogDeleteExecutorService.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + persistLogDeleteExecutorService, this::checkDeletePersistRaftLog, LOG_DELETE_CHECK_INTERVAL_SECOND, LOG_DELETE_CHECK_INTERVAL_SECOND, diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/server/PullSnapshotHintService.java b/cluster/src/main/java/org/apache/iotdb/cluster/server/PullSnapshotHintService.java index 3cec37e553..ad2df7582e 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/server/PullSnapshotHintService.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/server/PullSnapshotHintService.java @@ -29,6 +29,7 @@ import org.apache.iotdb.cluster.rpc.thrift.Node; import org.apache.iotdb.cluster.rpc.thrift.RaftNode; import org.apache.iotdb.cluster.server.member.DataGroupMember; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.thrift.TException; import org.slf4j.Logger; @@ -57,7 +58,8 @@ public class PullSnapshotHintService { public void start() { this.service = IoTDBThreadPoolFactory.newScheduledThreadPool(1, "PullSnapshotHint"); - this.service.scheduleAtFixedRate(this::sendHints, 0, 10, TimeUnit.MILLISECONDS); + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + this.service, this::sendHints, 0, 10, TimeUnit.MILLISECONDS); } public void stop() { diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/utils/nodetool/ClusterMonitor.java b/cluster/src/main/java/org/apache/iotdb/cluster/utils/nodetool/ClusterMonitor.java index 47d0232468..9763834128 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/utils/nodetool/ClusterMonitor.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/utils/nodetool/ClusterMonitor.java @@ -35,6 +35,7 @@ import org.apache.iotdb.cluster.utils.ClientUtils; import org.apache.iotdb.cluster.utils.nodetool.function.NodeToolCmd; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.exception.MetadataException; import org.apache.iotdb.commons.exception.StartupException; @@ -92,19 +93,20 @@ public class ClusterMonitor implements ClusterMonitorMBean, IService { private void startCollectClusterStatus() { // monitor all nodes' live status LOGGER.info("start metric node status and leader distribution"); - IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.CLUSTER_MONITOR.getName()) - .scheduleAtFixedRate( - () -> { - MetaGroupMember metaGroupMember = ClusterIoTDB.getInstance().getMetaGroupMember(); - if (metaGroupMember != null - && metaGroupMember.getLeader().equals(metaGroupMember.getThisNode())) { - metricNodeStatus(metaGroupMember); - metricLeaderDistribution(metaGroupMember); - } - }, - 10L, - 10L, - TimeUnit.SECONDS); + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( + ThreadName.CLUSTER_MONITOR.getName()), + () -> { + MetaGroupMember metaGroupMember = ClusterIoTDB.getInstance().getMetaGroupMember(); + if (metaGroupMember != null + && metaGroupMember.getLeader().equals(metaGroupMember.getThisNode())) { + metricNodeStatus(metaGroupMember); + metricLeaderDistribution(metaGroupMember); + } + }, + 10L, + 10L, + TimeUnit.SECONDS); } private void metricLeaderDistribution(MetaGroupMember metaGroupMember) { diff --git a/confignode/src/main/java/org/apache/iotdb/confignode/manager/PartitionManager.java b/confignode/src/main/java/org/apache/iotdb/confignode/manager/PartitionManager.java index d6e77c4de2..541573b301 100644 --- a/confignode/src/main/java/org/apache/iotdb/confignode/manager/PartitionManager.java +++ b/confignode/src/main/java/org/apache/iotdb/confignode/manager/PartitionManager.java @@ -25,6 +25,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.common.rpc.thrift.TSeriesPartitionSlot; import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.exception.MetadataException; import org.apache.iotdb.commons.partition.executor.SeriesPartitionExecutor; import org.apache.iotdb.confignode.client.SyncDataNodeClientPool; @@ -82,7 +83,8 @@ public class PartitionManager { this.partitionInfo = partitionInfo; this.regionCleaner = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("IoTDB-Region-Cleaner"); - regionCleaner.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + regionCleaner, this::clearDeletedRegions, REGION_CLEANER_WORK_INITIAL_DELAY, REGION_CLEANER_WORK_INTERVAL, diff --git a/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/IoTDBSink.java b/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/IoTDBSink.java index 8ab06c0975..1ccf737956 100644 --- a/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/IoTDBSink.java +++ b/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/IoTDBSink.java @@ -86,6 +86,7 @@ public class IoTDBSink<IN> extends RichSinkFunction<IN> { sessionPoolSize); } + @SuppressWarnings("unsafeThreadSchedule") void initScheduler() { if (batchSize > 0) { scheduledExecutor = Executors.newSingleThreadScheduledExecutor(); diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/ScheduledExecutorUtil.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/ScheduledExecutorUtil.java new file mode 100644 index 0000000000..fbf8bbf8ba --- /dev/null +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/ScheduledExecutorUtil.java @@ -0,0 +1,106 @@ +/* + * 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.commons.concurrent.threadpool; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; + +public class ScheduledExecutorUtil { + + private static final Logger logger = LoggerFactory.getLogger(ScheduledExecutorUtil.class); + + /** + * A safe wrapper method to make sure the exception thrown by the previous running will not affect + * the next one. Please reference the javadoc of {@link + * ScheduledExecutorService#scheduleAtFixedRate(Runnable, long, long, TimeUnit)} for more details. + * + * @param executor the ScheduledExecutorService instance. + * @param command same parameter in {@link ScheduledExecutorService#scheduleAtFixedRate(Runnable, + * long, long, TimeUnit)}. + * @param initialDelay same parameter in {@link + * ScheduledExecutorService#scheduleAtFixedRate(Runnable, long, long, TimeUnit)}. + * @param period same parameter in {@link ScheduledExecutorService#scheduleAtFixedRate(Runnable, + * long, long, TimeUnit)}. + * @param unit same parameter in {@link ScheduledExecutorService#scheduleAtFixedRate(Runnable, + * long, long, TimeUnit)}. + * @return the same return value of {@link ScheduledExecutorService#scheduleAtFixedRate(Runnable, + * long, long, TimeUnit)}. + */ + @SuppressWarnings("unsafeThreadSchedule") + public static ScheduledFuture<?> safelyScheduleAtFixedRate( + ScheduledExecutorService executor, + Runnable command, + long initialDelay, + long period, + TimeUnit unit) { + return executor.scheduleAtFixedRate( + () -> { + try { + command.run(); + } catch (Throwable t) { + logger.error("Schedule task failed", t); + } + }, + initialDelay, + period, + unit); + } + + /** + * A safe wrapper method to make sure the exception thrown by the previous running will not affect + * the next one. Please reference the javadoc of {@link + * ScheduledExecutorService#scheduleWithFixedDelay(Runnable, long, long, TimeUnit)} for more + * details. + * + * @param executor the ScheduledExecutorService instance. + * @param command same parameter in {@link + * ScheduledExecutorService#scheduleWithFixedDelay(Runnable, long, long, TimeUnit)}. + * @param initialDelay same parameter in {@link + * ScheduledExecutorService#scheduleWithFixedDelay(Runnable, long, long, TimeUnit)}. + * @param delay same parameter in {@link ScheduledExecutorService#scheduleWithFixedDelay(Runnable, + * long, long, TimeUnit)}. + * @param unit same parameter in {@link ScheduledExecutorService#scheduleWithFixedDelay(Runnable, + * long, long, TimeUnit)}. + * @return the same return value of {@link + * ScheduledExecutorService#scheduleWithFixedDelay(Runnable, long, long, TimeUnit)}. + */ + @SuppressWarnings("unsafeThreadSchedule") + public static ScheduledFuture<?> safelyScheduleWithFixedDelay( + ScheduledExecutorService executor, + Runnable command, + long initialDelay, + long delay, + TimeUnit unit) { + return executor.scheduleWithFixedDelay( + () -> { + try { + command.run(); + } catch (Throwable t) { + logger.error("Schedule task failed", t); + } + }, + initialDelay, + delay, + unit); + } +} diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedScheduledExecutorService.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedScheduledExecutorService.java index 74c15e43f5..a362983785 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedScheduledExecutorService.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedScheduledExecutorService.java @@ -58,12 +58,14 @@ public class WrappedScheduledExecutorService } @Override + @SuppressWarnings("unsafeThreadSchedule") public ScheduledFuture<?> scheduleAtFixedRate( Runnable command, long initialDelay, long period, TimeUnit unit) { return service.scheduleAtFixedRate(command, initialDelay, period, unit); } @Override + @SuppressWarnings("unsafeThreadSchedule") public ScheduledFuture<?> scheduleWithFixedDelay( Runnable command, long initialDelay, long delay, TimeUnit unit) { return service.scheduleWithFixedDelay(command, initialDelay, delay, unit); diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadScheduledExecutor.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadScheduledExecutor.java index 8410eebcbc..40d5e7183c 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadScheduledExecutor.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadScheduledExecutor.java @@ -56,12 +56,14 @@ public class WrappedSingleThreadScheduledExecutor } @Override + @SuppressWarnings("unsafeThreadSchedule") public ScheduledFuture<?> scheduleAtFixedRate( Runnable command, long initialDelay, long period, TimeUnit unit) { return service.scheduleAtFixedRate(command, initialDelay, period, unit); } @Override + @SuppressWarnings("unsafeThreadSchedule") public ScheduledFuture<?> scheduleWithFixedDelay( Runnable command, long initialDelay, long delay, TimeUnit unit) { return service.scheduleWithFixedDelay(command, initialDelay, delay, unit); diff --git a/node-commons/src/test/java/org/apache/iotdb/commons/IoTDBThreadPoolFactoryTest.java b/node-commons/src/test/java/org/apache/iotdb/commons/IoTDBThreadPoolFactoryTest.java index a58f274d48..f4af366000 100644 --- a/node-commons/src/test/java/org/apache/iotdb/commons/IoTDBThreadPoolFactoryTest.java +++ b/node-commons/src/test/java/org/apache/iotdb/commons/IoTDBThreadPoolFactoryTest.java @@ -20,6 +20,7 @@ package org.apache.iotdb.commons; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.WrappedRunnable; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.thrift.server.TThreadPoolServer; import org.apache.thrift.server.TThreadPoolServer.Args; @@ -120,7 +121,8 @@ public class IoTDBThreadPoolFactoryTest { IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(POOL_NAME, handler); for (int i = 0; i < threadCount; i++) { Runnable task = new TestThread(reason); - ScheduledFuture<?> future = exec.scheduleAtFixedRate(task, 0, 1, TimeUnit.SECONDS); + ScheduledFuture<?> future = + ScheduledExecutorUtil.safelyScheduleAtFixedRate(exec, task, 0, 1, TimeUnit.SECONDS); try { future.get(); } catch (ExecutionException e) { @@ -147,7 +149,8 @@ public class IoTDBThreadPoolFactoryTest { IoTDBThreadPoolFactory.newScheduledThreadPool(threadCount / 2, POOL_NAME, handler); for (int i = 0; i < threadCount; i++) { Runnable task = new TestThread(reason); - ScheduledFuture<?> future = exec.scheduleAtFixedRate(task, 0, 1, TimeUnit.SECONDS); + ScheduledFuture<?> future = + ScheduledExecutorUtil.safelyScheduleAtFixedRate(exec, task, 0, 1, TimeUnit.SECONDS); try { future.get(); } catch (ExecutionException e) { diff --git a/pom.xml b/pom.xml index 41cf1cc480..9dba3bff75 100644 --- a/pom.xml +++ b/pom.xml @@ -942,12 +942,12 @@ <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-checkstyle-plugin</artifactId> - <version>3.0.0</version> + <version>3.1.2</version> <dependencies> <dependency> <groupId>com.puppycrawl.tools</groupId> <artifactId>checkstyle</artifactId> - <version>8.18</version> + <version>8.45.1</version> </dependency> </dependencies> <executions> 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 a91ef781ff..aa94ff8a43 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 @@ -21,6 +21,7 @@ package org.apache.iotdb.db.engine; import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.commons.exception.MetadataException; import org.apache.iotdb.commons.exception.ShutdownException; @@ -289,8 +290,12 @@ public class StorageEngine implements IService { recover(); ttlCheckThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("TTL-Check"); - ttlCheckThread.scheduleAtFixedRate( - this::checkTTL, TTL_CHECK_INTERVAL, TTL_CHECK_INTERVAL, TimeUnit.MILLISECONDS); + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + ttlCheckThread, + this::checkTTL, + TTL_CHECK_INTERVAL, + TTL_CHECK_INTERVAL, + TimeUnit.MILLISECONDS); logger.info("start ttl check thread successfully."); startTimedService(); @@ -314,7 +319,8 @@ public class StorageEngine implements IService { seqMemtableTimedFlushCheckThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( ThreadName.TIMED_FlUSH_SEQ_MEMTABLE.getName()); - seqMemtableTimedFlushCheckThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + seqMemtableTimedFlushCheckThread, this::timedFlushSeqMemTable, config.getSeqMemtableFlushCheckInterval(), config.getSeqMemtableFlushCheckInterval(), @@ -326,7 +332,8 @@ public class StorageEngine implements IService { unseqMemtableTimedFlushCheckThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( ThreadName.TIMED_FlUSH_UNSEQ_MEMTABLE.getName()); - unseqMemtableTimedFlushCheckThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + unseqMemtableTimedFlushCheckThread, this::timedFlushUnseqMemTable, config.getUnseqMemtableFlushCheckInterval(), config.getUnseqMemtableFlushCheckInterval(), @@ -927,7 +934,11 @@ public class StorageEngine implements IService { } } - /** @return TsFiles (seq or unseq) grouped by their storage group and partition number. */ + /** + * Get all the closed tsfiles of each storage group. + * + * @return TsFiles (seq or unseq) grouped by their storage group and partition number. + */ public Map<PartialPath, Map<Long, List<TsFileResource>>> getAllClosedStorageGroupTsFile() { Map<PartialPath, Map<Long, List<TsFileResource>>> ret = new HashMap<>(); for (Entry<PartialPath, StorageGroupManager> entry : processorMap.entrySet()) { @@ -1031,7 +1042,11 @@ public class StorageEngine implements IService { list.forEach(DataRegion::readUnlock); } - /** @return virtual storage group name, like root.sg1/0 */ + /** + * Get the virtual storage group name. + * + * @return virtual storage group name, like root.sg1/0 + */ public String getStorageGroupPath(PartialPath path) throws StorageEngineException { PartialPath deviceId = path.getDevicePath(); DataRegion storageGroupProcessor = getProcessor(deviceId); diff --git a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngineV2.java b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngineV2.java index 7f7aa5ee80..f54451c9be 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngineV2.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngineV2.java @@ -22,6 +22,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.consensus.DataRegionId; import org.apache.iotdb.commons.exception.MetadataException; import org.apache.iotdb.commons.exception.ShutdownException; @@ -325,8 +326,12 @@ public class StorageEngineV2 implements IService { recover(); ttlCheckThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("TTL-Check"); - ttlCheckThread.scheduleAtFixedRate( - this::checkTTL, TTL_CHECK_INTERVAL, TTL_CHECK_INTERVAL, TimeUnit.MILLISECONDS); + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + ttlCheckThread, + this::checkTTL, + TTL_CHECK_INTERVAL, + TTL_CHECK_INTERVAL, + TimeUnit.MILLISECONDS); logger.info("start ttl check thread successfully."); startTimedService(); @@ -352,7 +357,8 @@ public class StorageEngineV2 implements IService { seqMemtableTimedFlushCheckThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( ThreadName.TIMED_FlUSH_SEQ_MEMTABLE.getName()); - seqMemtableTimedFlushCheckThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + seqMemtableTimedFlushCheckThread, this::timedFlushSeqMemTable, config.getSeqMemtableFlushCheckInterval(), config.getSeqMemtableFlushCheckInterval(), @@ -364,7 +370,8 @@ public class StorageEngineV2 implements IService { unseqMemtableTimedFlushCheckThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( ThreadName.TIMED_FlUSH_UNSEQ_MEMTABLE.getName()); - unseqMemtableTimedFlushCheckThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + unseqMemtableTimedFlushCheckThread, this::timedFlushUnseqMemTable, config.getUnseqMemtableFlushCheckInterval(), config.getUnseqMemtableFlushCheckInterval(), diff --git a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManager.java b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManager.java index 500434d9a2..4b07aa124a 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManager.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManager.java @@ -21,6 +21,7 @@ package org.apache.iotdb.db.engine.compaction; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.concurrent.threadpool.WrappedScheduledExecutorService; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.service.IService; @@ -115,14 +116,9 @@ public class CompactionTaskManager implements IService { // candidateCompactionTaskQueue, check that all tsfiles in the compaction task are valid, and // if there is thread space available in the taskExecutionPool, put the compaction task thread // into the taskExecutionPool and perform the compaction. - compactionTaskSubmissionThreadPool.scheduleWithFixedDelay( - () -> { - try { - submitTaskFromTaskQueue(); - } catch (Throwable t) { - logger.error("Schedule {} failed", ThreadName.COMPACTION_SERVICE.getName(), t); - } - }, + ScheduledExecutorUtil.safelyScheduleWithFixedDelay( + compactionTaskSubmissionThreadPool, + this::submitTaskFromTaskQueue, TASK_SUBMIT_INTERVAL, TASK_SUBMIT_INTERVAL, TimeUnit.MILLISECONDS); diff --git a/server/src/main/java/org/apache/iotdb/db/engine/cq/ContinuousQueryService.java b/server/src/main/java/org/apache/iotdb/db/engine/cq/ContinuousQueryService.java index d4882a358d..e7f2e7275c 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/cq/ContinuousQueryService.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/cq/ContinuousQueryService.java @@ -20,6 +20,7 @@ package org.apache.iotdb.db.engine.cq; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.exception.StartupException; import org.apache.iotdb.commons.file.SystemFileFactory; import org.apache.iotdb.commons.service.IService; @@ -126,7 +127,8 @@ public class ContinuousQueryService implements IService { continuousQueryTaskSubmitThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("CQ-Task-Submit-Thread"); - continuousQueryTaskSubmitThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + continuousQueryTaskSubmitThread, this::checkAndSubmitTasks, 0, TASK_SUBMIT_CHECK_INTERVAL, diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java index 2fc87803fe..480311094a 100755 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java @@ -21,6 +21,7 @@ package org.apache.iotdb.db.engine.storagegroup; import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.commons.exception.MetadataException; @@ -520,14 +521,9 @@ public class DataRegion { + logicalStorageGroupName + "-" + dataRegionId); - timedCompactionScheduleTask.scheduleWithFixedDelay( - () -> { - try { - executeCompaction(); - } catch (Throwable t) { - logger.error("Schedule {} failed", ThreadName.COMPACTION_SCHEDULE.getName(), t); - } - }, + ScheduledExecutorUtil.safelyScheduleWithFixedDelay( + timedCompactionScheduleTask, + this::executeCompaction, COMPACTION_TASK_SUBMIT_DELAY, IoTDBDescriptor.getInstance().getConfig().getCompactionScheduleIntervalInMs(), TimeUnit.MILLISECONDS); @@ -1109,7 +1105,11 @@ public class DataRegion { } } - /** @return whether the given time falls in ttl */ + /** + * Check whether the time falls in TTL. + * + * @return whether the given time falls in ttl + */ private boolean isAlive(long time) { return dataTTL == Long.MAX_VALUE || (System.currentTimeMillis() - time) <= dataTTL; } @@ -3178,7 +3178,11 @@ public class DataRegion { return dataRegionId; } - /** @return data region path, like root.sg1/0 */ + /** + * Get the storageGroupPath with dataRegionId. + * + * @return data region path, like root.sg1/0 + */ public String getStorageGroupPath() { return logicalStorageGroupName + File.separator + dataRegionId; } diff --git a/server/src/main/java/org/apache/iotdb/db/localconfignode/LocalConfigNode.java b/server/src/main/java/org/apache/iotdb/db/localconfignode/LocalConfigNode.java index a48887b687..0951b00525 100644 --- a/server/src/main/java/org/apache/iotdb/db/localconfignode/LocalConfigNode.java +++ b/server/src/main/java/org/apache/iotdb/db/localconfignode/LocalConfigNode.java @@ -20,6 +20,7 @@ package org.apache.iotdb.db.localconfignode; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.consensus.DataRegionId; import org.apache.iotdb.commons.consensus.SchemaRegionId; @@ -143,8 +144,8 @@ public class LocalConfigNode { if (config.getSyncMlogPeriodInMs() != 0) { timedForceMLogThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("timedForceMLogThread"); - - timedForceMLogThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + timedForceMLogThread, this::forceMlog, config.getSyncMlogPeriodInMs(), config.getSyncMlogPeriodInMs(), diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/fragment/FragmentInstanceManager.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/fragment/FragmentInstanceManager.java index 29d41fd934..e7099c3f75 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/fragment/FragmentInstanceManager.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/fragment/FragmentInstanceManager.java @@ -19,6 +19,7 @@ package org.apache.iotdb.db.mpp.execution.fragment; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.db.engine.storagegroup.DataRegion; import org.apache.iotdb.db.metadata.schemaregion.ISchemaRegion; import org.apache.iotdb.db.mpp.common.FragmentInstanceId; @@ -75,17 +76,8 @@ public class FragmentInstanceManager { this.infoCacheTime = new Duration(15, TimeUnit.MINUTES); - instanceManagementExecutor.scheduleWithFixedDelay( - () -> { - try { - removeOldInstances(); - } catch (Throwable e) { - logger.warn("Error removing old tasks", e); - } - }, - 200, - 200, - TimeUnit.MILLISECONDS); + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + instanceManagementExecutor, this::removeOldInstances, 200, 200, TimeUnit.MILLISECONDS); } public FragmentInstanceInfo execDataQueryFragmentInstance( diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FixedRateFragInsStateTracker.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FixedRateFragInsStateTracker.java index ac748bf2a8..755b628e91 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FixedRateFragInsStateTracker.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FixedRateFragInsStateTracker.java @@ -22,6 +22,7 @@ package org.apache.iotdb.db.mpp.plan.scheduler; import org.apache.iotdb.common.rpc.thrift.TEndPoint; import org.apache.iotdb.commons.client.IClientManager; import org.apache.iotdb.commons.client.sync.SyncDataNodeInternalServiceClient; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.db.mpp.execution.QueryStateMachine; import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceState; import org.apache.iotdb.db.mpp.plan.planner.plan.FragmentInstance; @@ -57,8 +58,12 @@ public class FixedRateFragInsStateTracker extends AbstractFragInsStateTracker { @Override public void start() { trackTask = - scheduledExecutor.scheduleAtFixedRate( - this::fetchStateAndUpdate, 0, STATE_FETCH_INTERVAL_IN_MS, TimeUnit.MILLISECONDS); + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + scheduledExecutor, + this::fetchStateAndUpdate, + 0, + STATE_FETCH_INTERVAL_IN_MS, + TimeUnit.MILLISECONDS); } @Override diff --git a/server/src/main/java/org/apache/iotdb/db/query/control/SessionTimeoutManager.java b/server/src/main/java/org/apache/iotdb/db/query/control/SessionTimeoutManager.java index 6959aa9a55..bab759a836 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/control/SessionTimeoutManager.java +++ b/server/src/main/java/org/apache/iotdb/db/query/control/SessionTimeoutManager.java @@ -19,6 +19,7 @@ package org.apache.iotdb.db.query.control; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.service.basic.ServiceProvider; @@ -49,7 +50,8 @@ public class SessionTimeoutManager { this.executorService = IoTDBThreadPoolFactory.newScheduledThreadPool(1, "session-timeout-manager"); - executorService.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + executorService, () -> { LOGGER.info("cleaning up expired sessions"); cleanup(); diff --git a/server/src/main/java/org/apache/iotdb/db/service/basic/QueryFrequencyRecorder.java b/server/src/main/java/org/apache/iotdb/db/service/basic/QueryFrequencyRecorder.java index f69d274958..5c4290f4c4 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/basic/QueryFrequencyRecorder.java +++ b/server/src/main/java/org/apache/iotdb/db/service/basic/QueryFrequencyRecorder.java @@ -19,6 +19,7 @@ package org.apache.iotdb.db.service.basic; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.db.conf.IoTDBConfig; import org.slf4j.Logger; @@ -36,7 +37,8 @@ public class QueryFrequencyRecorder { public QueryFrequencyRecorder(IoTDBConfig config) { ScheduledExecutorService timedQuerySqlCountThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor("timedQuerySqlCount"); - timedQuerySqlCountThread.scheduleAtFixedRate( + ScheduledExecutorUtil.safelyScheduleAtFixedRate( + timedQuerySqlCountThread, () -> { if (QUERY_COUNT.get() != 0) { QUERY_FREQUENCY_LOGGER.info( diff --git a/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java b/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java index 92818432f7..720557f235 100644 --- a/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java +++ b/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java @@ -21,6 +21,7 @@ package org.apache.iotdb.db.sync.sender.service; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.utils.TestOnly; import org.apache.iotdb.db.exception.SyncConnectionException; import org.apache.iotdb.db.sync.conf.SyncConstant; @@ -95,17 +96,9 @@ public class TransportHandler { public void start() { transportFuture = transportExecutorService.submit(transportClient); heartbeatFuture = - heartbeatExecutorService.scheduleWithFixedDelay( - () -> { - try { - sendHeartbeat(); - } catch (Throwable t) { - logger.error( - "Schedule {} failed", - ThreadName.SYNC_SENDER_HEARTBEAT.getName() + "-" + pipeName, - t); - } - }, + ScheduledExecutorUtil.safelyScheduleWithFixedDelay( + heartbeatExecutorService, + this::sendHeartbeat, 0, SyncConstant.HEARTBEAT_DELAY_SECONDS, TimeUnit.SECONDS); diff --git a/server/src/main/java/org/apache/iotdb/db/wal/WALManager.java b/server/src/main/java/org/apache/iotdb/db/wal/WALManager.java index 57d3dca4d6..5fd4de0e27 100644 --- a/server/src/main/java/org/apache/iotdb/db/wal/WALManager.java +++ b/server/src/main/java/org/apache/iotdb/db/wal/WALManager.java @@ -20,6 +20,7 @@ package org.apache.iotdb.db.wal; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.exception.StartupException; import org.apache.iotdb.commons.service.IService; import org.apache.iotdb.commons.service.ServiceType; @@ -175,17 +176,8 @@ public class WALManager implements IService { private void registerScheduleTask(long initDelayMs, long periodMs) { walDeleteThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.WAL_DELETE.getName()); - walDeleteThread.scheduleWithFixedDelay( - () -> { - try { - deleteOutdatedFiles(); - } catch (Throwable t) { - logger.error("Schedule {} failed", ThreadName.WAL_DELETE.getName(), t); - } - }, - initDelayMs, - periodMs, - TimeUnit.MILLISECONDS); + ScheduledExecutorUtil.safelyScheduleWithFixedDelay( + walDeleteThread, this::deleteOutdatedFiles, initDelayMs, periodMs, TimeUnit.MILLISECONDS); } @TestOnly
