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