This is an automated email from the ASF dual-hosted git repository.

qiaojialin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 5cd0a21f94 [IOTDB-3402] Fix abuse of 
ScheduledExecutorService.scheduleWithFixedDelay (#6172)
5cd0a21f94 is described below

commit 5cd0a21f94603939ac0d67ecc3e2c688250d97e3
Author: BaiJian <[email protected]>
AuthorDate: Mon Jun 6 22:55:33 2022 +0800

    [IOTDB-3402] Fix abuse of ScheduledExecutorService.scheduleWithFixedDelay 
(#6172)
---
 .../engine/compaction/CompactionTaskManager.java   |  8 +++++-
 .../iotdb/db/engine/storagegroup/DataRegion.java   |  8 +++++-
 .../db/sync/sender/service/TransportHandler.java   | 14 +++++++++-
 .../java/org/apache/iotdb/db/wal/WALManager.java   | 30 ++++++++++++++--------
 4 files changed, 46 insertions(+), 14 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..500434d9a2 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
@@ -116,7 +116,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..2fc87803fe 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
@@ -521,7 +521,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);
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..92818432f7 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
@@ -96,7 +96,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..57d3dca4d6 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
@@ -96,13 +96,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 +114,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 +172,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