This is an automated email from the ASF dual-hosted git repository. ericpai pushed a commit to branch bugfix/iotdb-3402 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 5b12e9ed3ba08d9c6ce5a91f895f1bdadd463cf9 Author: ericpai <[email protected]> AuthorDate: Mon Jun 6 18:01:11 2022 +0800 [IOTDB-3402] Fix abuse of ScheduledExecutorService.scheduleWithFixedDelay --- .../engine/compaction/CompactionTaskManager.java | 18 ++++---- .../iotdb/db/engine/storagegroup/DataRegion.java | 48 ++++++++-------------- .../db/sync/sender/service/TransportHandler.java | 15 ++++++- .../java/org/apache/iotdb/db/wal/WALManager.java | 37 +++++++++-------- 4 files changed, 59 insertions(+), 59 deletions(-) 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 385e23aa0d..6956b1169e 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 @@ -19,6 +19,7 @@ package org.apache.iotdb.db.engine.compaction; +import com.google.common.util.concurrent.RateLimiter; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; import org.apache.iotdb.commons.concurrent.threadpool.WrappedScheduledExecutorService; @@ -32,8 +33,6 @@ import org.apache.iotdb.db.engine.compaction.constant.CompactionTaskStatus; import org.apache.iotdb.db.engine.compaction.task.AbstractCompactionTask; import org.apache.iotdb.db.engine.compaction.task.CompactionTaskSummary; import org.apache.iotdb.db.utils.datastructure.FixedPriorityBlockingQueue; - -import com.google.common.util.concurrent.RateLimiter; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -41,12 +40,7 @@ import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.concurrent.Callable; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Future; -import java.util.concurrent.RejectedExecutionException; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; +import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; /** CompactionMergeTaskPoolManager provides a ThreadPool tPro queue and run all compaction tasks. */ @@ -116,7 +110,13 @@ public class CompactionTaskManager implements IService { // if there is thread space available in the taskExecutionPool, put the compaction task thread // into the taskExecutionPool and perform the compaction. compactionTaskSubmissionThreadPool.scheduleWithFixedDelay( - this::submitTaskFromTaskQueue, + () -> { + try { + submitTaskFromTaskQueue(); + } catch (Throwable t) { + logger.error("Schedule {} failed", ThreadName.COMPACTION_SERVICE.getName(), t); + } + }, TASK_SUBMIT_INTERVAL, TASK_SUBMIT_INTERVAL, TimeUnit.MILLISECONDS); 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 8f5610194e..3d0a511940 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 @@ -18,6 +18,7 @@ */ package org.apache.iotdb.db.engine.storagegroup; +import org.apache.commons.io.FileUtils; import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; @@ -48,26 +49,14 @@ import org.apache.iotdb.db.engine.upgrade.UpgradeCheckStatus; import org.apache.iotdb.db.engine.upgrade.UpgradeLog; import org.apache.iotdb.db.engine.version.SimpleFileVersionController; import org.apache.iotdb.db.engine.version.VersionController; -import org.apache.iotdb.db.exception.BatchProcessException; -import org.apache.iotdb.db.exception.DataRegionException; -import org.apache.iotdb.db.exception.DiskSpaceInsufficientException; -import org.apache.iotdb.db.exception.LoadFileException; -import org.apache.iotdb.db.exception.TriggerExecutionException; -import org.apache.iotdb.db.exception.TsFileProcessorException; -import org.apache.iotdb.db.exception.WriteProcessException; -import org.apache.iotdb.db.exception.WriteProcessRejectException; +import org.apache.iotdb.db.exception.*; import org.apache.iotdb.db.exception.query.OutOfTTLException; import org.apache.iotdb.db.exception.query.QueryProcessException; import org.apache.iotdb.db.metadata.cache.DataNodeSchemaCache; import org.apache.iotdb.db.metadata.idtable.IDTable; import org.apache.iotdb.db.metadata.idtable.IDTableManager; import org.apache.iotdb.db.metadata.mnode.IMeasurementMNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertMultiTabletsNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertRowNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertRowsNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertRowsOfOneDeviceNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertTabletNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.*; import org.apache.iotdb.db.qp.physical.crud.DeletePlan; import org.apache.iotdb.db.qp.physical.crud.InsertRowPlan; import org.apache.iotdb.db.qp.physical.crud.InsertRowsOfOneDevicePlan; @@ -103,27 +92,14 @@ import org.apache.iotdb.tsfile.fileSystem.fsFactory.FSFactory; import org.apache.iotdb.tsfile.read.filter.basic.Filter; import org.apache.iotdb.tsfile.utils.Pair; import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter; - -import org.apache.commons.io.FileUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.File; import java.io.IOException; import java.nio.file.Files; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.Collection; -import java.util.Collections; -import java.util.Date; -import java.util.HashMap; -import java.util.Iterator; -import java.util.LinkedList; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.Map.Entry; -import java.util.Set; -import java.util.TreeMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -521,7 +497,13 @@ public class DataRegion { + "-" + dataRegionId); timedCompactionScheduleTask.scheduleWithFixedDelay( - this::executeCompaction, + () -> { + try { + executeCompaction(); + } catch (Throwable t) { + logger.error("Schedule {} failed", ThreadName.COMPACTION_SCHEDULE.getName(), t); + } + }, COMPACTION_TASK_SUBMIT_DELAY, IoTDBDescriptor.getInstance().getConfig().getCompactionScheduleIntervalInMs(), TimeUnit.MILLISECONDS); @@ -1103,7 +1085,9 @@ public class DataRegion { } } - /** @return whether the given 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; } @@ -3172,7 +3156,9 @@ public class DataRegion { return dataRegionId; } - /** @return data region path, like root.sg1/0 */ + /** + * @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/sync/sender/service/TransportHandler.java b/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java index af49ab9d2e..a3846e8648 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 @@ -31,7 +31,6 @@ import org.apache.iotdb.db.sync.transport.client.TransportClient; import org.apache.iotdb.service.transport.thrift.RequestType; import org.apache.iotdb.service.transport.thrift.SyncRequest; import org.apache.iotdb.service.transport.thrift.SyncResponse; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -96,7 +95,19 @@ public class TransportHandler { transportFuture = transportExecutorService.submit(transportClient); heartbeatFuture = heartbeatExecutorService.scheduleWithFixedDelay( - this::sendHeartbeat, 0, SyncConstant.HEARTBEAT_DELAY_SECONDS, TimeUnit.SECONDS); + () -> { + try { + sendHeartbeat(); + } catch (Throwable t) { + logger.error( + "Schedule {} failed", + ThreadName.SYNC_SENDER_HEARTBEAT.getName() + "-" + pipeName, + t); + } + }, + 0, + SyncConstant.HEARTBEAT_DELAY_SECONDS, + TimeUnit.SECONDS); } public void stop() { 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 8e028ad09a..bcdde79533 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 @@ -34,15 +34,10 @@ import org.apache.iotdb.db.wal.node.IWALNode; import org.apache.iotdb.db.wal.node.WALFakeNode; import org.apache.iotdb.db.wal.node.WALNode; import org.apache.iotdb.db.wal.utils.WALMode; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Future; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; +import java.util.concurrent.*; /** This class is used to manage and allocate wal nodes */ public class WALManager implements IService { @@ -96,13 +91,8 @@ public class WALManager implements IService { } try { - walDeleteThread = - IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.WAL_DELETE.getName()); - walDeleteThread.scheduleWithFixedDelay( - this::deleteOutdatedFiles, - config.getDeleteWalFilesPeriodInMs(), - config.getDeleteWalFilesPeriodInMs(), - TimeUnit.MILLISECONDS); + registerScheduleTask( + config.getDeleteWalFilesPeriodInMs(), config.getDeleteWalFilesPeriodInMs()); } catch (Exception e) { throw new StartupException(this.getID().getName(), e.getMessage()); } @@ -119,10 +109,7 @@ public class WALManager implements IService { shutdownThread(walDeleteThread, ThreadName.WAL_DELETE); } logger.info("Stop wal delete thread successfully, and now restart it."); - walDeleteThread = - IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.WAL_DELETE.getName()); - walDeleteThread.scheduleWithFixedDelay( - this::deleteOutdatedFiles, 0, config.getDeleteWalFilesPeriodInMs(), TimeUnit.MILLISECONDS); + registerScheduleTask(0, config.getDeleteWalFilesPeriodInMs()); logger.info( "Reboot wal delete thread successfully, current period is {} ms", config.getDeleteWalFilesPeriodInMs()); @@ -180,6 +167,22 @@ 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); + } + @TestOnly public void clear() { walNodesManager.clear();
