You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@iotdb.apache.org by qi...@apache.org on 2021/07/13 06:13:19 UTC
[iotdb] branch rel/0.12 updated: [To rel/0.12] fix compaction block
flush bug (#3534)
This is an automated email from the ASF dual-hosted git repository.
qiaojialin pushed a commit to branch rel/0.12
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rel/0.12 by this push:
new bf6d83c [To rel/0.12] fix compaction block flush bug (#3534)
bf6d83c is described below
commit bf6d83cd45d195e4c9007e386e6b90aec2cf0fbb
Author: zhanglingzhe0820 <44...@qq.com>
AuthorDate: Tue Jul 13 14:12:57 2021 +0800
[To rel/0.12] fix compaction block flush bug (#3534)
---
.../db/engine/compaction/CompactionMergeTaskPoolManager.java | 11 +++++++++--
1 file changed, 9 insertions(+), 2 deletions(-)
diff --git a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionMergeTaskPoolManager.java b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionMergeTaskPoolManager.java
index 9f7ff5a..9b7949c 100644
--- a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionMergeTaskPoolManager.java
+++ b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionMergeTaskPoolManager.java
@@ -54,7 +54,8 @@ public class CompactionMergeTaskPoolManager implements IService {
LoggerFactory.getLogger(CompactionMergeTaskPoolManager.class);
private static final CompactionMergeTaskPoolManager INSTANCE =
new CompactionMergeTaskPoolManager();
- private ScheduledExecutorService pool;
+ private ScheduledExecutorService scheduledPool;
+ private ExecutorService pool;
private Map<String, Set<Future<Void>>> storageGroupTasks = new ConcurrentHashMap<>();
private static final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
@@ -71,6 +72,9 @@ public class CompactionMergeTaskPoolManager implements IService {
IoTDBThreadPoolFactory.newScheduledThreadPool(
IoTDBDescriptor.getInstance().getConfig().getCompactionThreadNum(),
ThreadName.COMPACTION_SERVICE.getName());
+ this.scheduledPool =
+ IoTDBThreadPoolFactory.newScheduledThreadPool(
+ Integer.MAX_VALUE, ThreadName.COMPACTION_SERVICE.getName());
}
logger.info("Compaction task manager started.");
}
@@ -78,6 +82,7 @@ public class CompactionMergeTaskPoolManager implements IService {
@Override
public void stop() {
if (pool != null) {
+ scheduledPool.shutdownNow();
pool.shutdownNow();
logger.info("Waiting for task pool to shut down");
waitTermination();
@@ -88,6 +93,7 @@ public class CompactionMergeTaskPoolManager implements IService {
@Override
public void waitAndStop(long milliseconds) {
if (pool != null) {
+ awaitTermination(scheduledPool, milliseconds);
awaitTermination(pool, milliseconds);
logger.info("Waiting for task pool to shut down");
waitTermination();
@@ -142,6 +148,7 @@ public class CompactionMergeTaskPoolManager implements IService {
logger.warn("CompactionManager has wait for {} seconds to stop", time / 1000);
}
}
+ scheduledPool = null;
pool = null;
storageGroupTasks.clear();
logger.info("CompactionManager stopped");
@@ -190,7 +197,7 @@ public class CompactionMergeTaskPoolManager implements IService {
}
public void init(Runnable function) {
- pool.scheduleWithFixedDelay(
+ scheduledPool.scheduleWithFixedDelay(
function, 1000, config.getCompactionInterval(), TimeUnit.MILLISECONDS);
}