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();

Reply via email to