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 58df85eee42567a5c3961fd5df0a6e61dc747c55 Author: ericpai <[email protected]> AuthorDate: Mon Jun 6 18:01:11 2022 +0800 [IOTDB-3402] Fix abuse of ScheduledExecutorService.scheduleWithFixedDelay --- .../engine/compaction/CompactionTaskManager.java | 11 +++++--- .../iotdb/db/engine/storagegroup/DataRegion.java | 19 +++++++++---- .../db/sync/sender/service/TransportHandler.java | 15 +++++++++-- .../java/org/apache/iotdb/db/wal/WALManager.java | 31 +++++++++++++--------- 4 files changed, 54 insertions(+), 22 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..5792b6fa0d 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; @@ -116,7 +115,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..d26aebfe4e 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; @@ -103,8 +104,6 @@ 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; @@ -521,7 +520,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 +1108,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 +3179,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..67ffba0e67 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,7 +34,6 @@ 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; @@ -96,13 +95,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 +113,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 +171,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();
