This is an automated email from the ASF dual-hosted git repository. rong pushed a commit to branch 0.12-revert in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 696a202f2b8be19ba3f5b6b981bac54047396e07 Author: Steve Yurong Su <[email protected]> AuthorDate: Tue Jul 27 16:20:48 2021 +0800 Revert "[To rel/0.12] fix compaction block flush bug (#3534)" This reverts commit bf6d83cd --- .../db/engine/compaction/CompactionMergeTaskPoolManager.java | 11 ++--------- 1 file changed, 2 insertions(+), 9 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 9b7949c..9f7ff5a 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,8 +54,7 @@ public class CompactionMergeTaskPoolManager implements IService { LoggerFactory.getLogger(CompactionMergeTaskPoolManager.class); private static final CompactionMergeTaskPoolManager INSTANCE = new CompactionMergeTaskPoolManager(); - private ScheduledExecutorService scheduledPool; - private ExecutorService pool; + private ScheduledExecutorService pool; private Map<String, Set<Future<Void>>> storageGroupTasks = new ConcurrentHashMap<>(); private static final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); @@ -72,9 +71,6 @@ 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."); } @@ -82,7 +78,6 @@ 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(); @@ -93,7 +88,6 @@ 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(); @@ -148,7 +142,6 @@ 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"); @@ -197,7 +190,7 @@ public class CompactionMergeTaskPoolManager implements IService { } public void init(Runnable function) { - scheduledPool.scheduleWithFixedDelay( + pool.scheduleWithFixedDelay( function, 1000, config.getCompactionInterval(), TimeUnit.MILLISECONDS); }
