This is an automated email from the ASF dual-hosted git repository. tanxinyu pushed a commit to branch master_performance in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit c0233e63077404607baa54e47f361f76ebefa5e9 Author: Xiangwei Wei <[email protected]> AuthorDate: Fri Nov 26 21:00:40 2021 +0800 [IOTDB-2065] TsFileSequenceReader will be cached for 100s even no longer used (#4478) --- .../iotdb/db/query/control/FileReaderManager.java | 110 +++++++-------------- .../iotdb/db/engine/cache/ChunkCacheTest.java | 1 - .../db/engine/compaction/cross/MergeTest.java | 1 - .../compaction/inner/InnerCompactionTest.java | 1 - .../SizeTieredCompactionRecoverTest.java | 1 - .../inner/sizetiered/SizeTieredCompactionTest.java | 1 - .../compaction/utils/CompactionClearUtils.java | 1 - .../db/query/control/FileReaderManagerTest.java | 19 +--- .../query/reader/series/SeriesReaderTestUtil.java | 1 - .../iotdb/db/rescon/ResourceManagerTest.java | 1 - 10 files changed, 37 insertions(+), 100 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/query/control/FileReaderManager.java b/server/src/main/java/org/apache/iotdb/db/query/control/FileReaderManager.java index 01cde88..0046bf0 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/control/FileReaderManager.java +++ b/server/src/main/java/org/apache/iotdb/db/query/control/FileReaderManager.java @@ -18,11 +18,7 @@ */ package org.apache.iotdb.db.query.control; -import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; -import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.engine.storagegroup.TsFileResource; -import org.apache.iotdb.db.service.IService; -import org.apache.iotdb.db.service.ServiceType; import org.apache.iotdb.tsfile.common.conf.TSFileConfig; import org.apache.iotdb.tsfile.read.TsFileSequenceReader; import org.apache.iotdb.tsfile.read.UnClosedTsFileReader; @@ -35,15 +31,13 @@ import java.io.IOException; import java.util.Iterator; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; /** * FileReaderManager is a singleton, which is used to manage all file readers(opened file streams) * to ensure that each file is opened at most once. */ -public class FileReaderManager implements IService { +public class FileReaderManager { private static final Logger logger = LoggerFactory.getLogger(FileReaderManager.class); private static final Logger resourceLogger = LoggerFactory.getLogger("FileMonitor"); @@ -73,16 +67,11 @@ public class FileReaderManager implements IService { */ private Map<String, AtomicInteger> unclosedReferenceMap; - private ScheduledExecutorService executorService; - private FileReaderManager() { closedFileReaderMap = new ConcurrentHashMap<>(); unclosedFileReaderMap = new ConcurrentHashMap<>(); closedReferenceMap = new ConcurrentHashMap<>(); unclosedReferenceMap = new ConcurrentHashMap<>(); - executorService = IoTDBThreadPoolFactory.newScheduledThreadPool(1, "open-files-manager"); - - clearUnUsedFilesInFixTime(); } public static FileReaderManager getInstance() { @@ -102,44 +91,6 @@ public class FileReaderManager implements IService { } } - private void clearUnUsedFilesInFixTime() { - long examinePeriod = IoTDBDescriptor.getInstance().getConfig().getCacheFileReaderClearPeriod(); - executorService.scheduleAtFixedRate( - () -> { - synchronized (this) { - clearMap(closedFileReaderMap, closedReferenceMap); - clearMap(unclosedFileReaderMap, unclosedReferenceMap); - } - }, - 0, - examinePeriod, - TimeUnit.MILLISECONDS); - } - - private void clearMap( - Map<String, TsFileSequenceReader> readerMap, Map<String, AtomicInteger> refMap) { - Iterator<Map.Entry<String, TsFileSequenceReader>> iterator = readerMap.entrySet().iterator(); - while (iterator.hasNext()) { - Map.Entry<String, TsFileSequenceReader> entry = iterator.next(); - TsFileSequenceReader reader = entry.getValue(); - AtomicInteger refAtom = refMap.get(entry.getKey()); - - if (refAtom != null && refAtom.get() == 0) { - try { - reader.close(); - } catch (IOException e) { - logger.error("Can not close TsFileSequenceReader {} !", reader.getFileName(), e); - } - iterator.remove(); - refMap.remove(entry.getKey()); - if (resourceLogger.isDebugEnabled()) { - resourceLogger.debug( - "{} TsFileReader is closed because of no reference.", entry.getKey()); - } - } - } - } - /** * Get the reader of the file(tsfile or unseq tsfile) indicated by filePath. If the reader already * exists, just get it from closedFileReaderMap or unclosedFileReaderMap depending on isClosing . @@ -211,14 +162,44 @@ public class FileReaderManager implements IService { void decreaseFileReaderReference(TsFileResource tsFile, boolean isClosed) { synchronized (this) { if (!isClosed && unclosedReferenceMap.containsKey(tsFile.getTsFilePath())) { - unclosedReferenceMap.get(tsFile.getTsFilePath()).decrementAndGet(); + if (unclosedReferenceMap.get(tsFile.getTsFilePath()).decrementAndGet() == 0) { + closeUnUsedReaderAndRemoveRef(tsFile.getTsFilePath(), false); + } } else if (closedReferenceMap.containsKey(tsFile.getTsFilePath())) { - closedReferenceMap.get(tsFile.getTsFilePath()).decrementAndGet(); + if (closedReferenceMap.get(tsFile.getTsFilePath()).decrementAndGet() == 0) { + closeUnUsedReaderAndRemoveRef(tsFile.getTsFilePath(), true); + } } } tsFile.readUnlock(); } + private void closeUnUsedReaderAndRemoveRef(String tsFilePath, boolean isClosed) { + Map<String, TsFileSequenceReader> readerMap = + isClosed ? closedFileReaderMap : unclosedFileReaderMap; + Map<String, AtomicInteger> refMap = isClosed ? closedReferenceMap : unclosedReferenceMap; + synchronized (this) { + // check ref num again + if (refMap.get(tsFilePath).get() != 0) { + return; + } + + TsFileSequenceReader reader = readerMap.get(tsFilePath); + if (reader != null) { + try { + reader.close(); + } catch (IOException e) { + logger.error("Can not close TsFileSequenceReader {} !", reader.getFileName(), e); + } + } + readerMap.remove(tsFilePath); + refMap.remove(tsFilePath); + if (resourceLogger.isDebugEnabled()) { + resourceLogger.debug("{} TsFileReader is closed because of no reference.", tsFilePath); + } + } + } + /** * Only for <code>EnvironmentUtils.cleanEnv</code> method. To make sure that unit tests and * integration tests will not conflict with each other. @@ -253,31 +234,6 @@ public class FileReaderManager implements IService { || (!isClosed && unclosedFileReaderMap.containsKey(tsFile.getTsFilePath())); } - @Override - public void start() { - // Do nothing - } - - @Override - public void stop() { - if (executorService == null || executorService.isShutdown()) { - return; - } - - executorService.shutdown(); - try { - executorService.awaitTermination(10, TimeUnit.SECONDS); - } catch (InterruptedException e) { - logger.error("StatMonitor timing service could not be shutdown.", e); - Thread.currentThread().interrupt(); - } - } - - @Override - public ServiceType getID() { - return ServiceType.FILE_READER_MANAGER_SERVICE; - } - private static class FileReaderManagerHelper { private static final FileReaderManager INSTANCE = new FileReaderManager(); diff --git a/server/src/test/java/org/apache/iotdb/db/engine/cache/ChunkCacheTest.java b/server/src/test/java/org/apache/iotdb/db/engine/cache/ChunkCacheTest.java index 33eab71..52fca3f 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/cache/ChunkCacheTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/cache/ChunkCacheTest.java @@ -234,6 +234,5 @@ public class ChunkCacheTest { resourceFile.delete(); } FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); - FileReaderManager.getInstance().stop(); } } diff --git a/server/src/test/java/org/apache/iotdb/db/engine/compaction/cross/MergeTest.java b/server/src/test/java/org/apache/iotdb/db/engine/compaction/cross/MergeTest.java index fe53e61..5882ddd 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/compaction/cross/MergeTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/compaction/cross/MergeTest.java @@ -173,7 +173,6 @@ abstract class MergeTest { } FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); - FileReaderManager.getInstance().stop(); } void prepareFile(TsFileResource tsFileResource, long timeOffset, long ptNum, long valueOffset) diff --git a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionTest.java b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionTest.java index 7559da5..c7aca33 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionTest.java @@ -205,7 +205,6 @@ public abstract class InnerCompactionTest { resourceFile.delete(); } FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); - FileReaderManager.getInstance().stop(); } void prepareFile(TsFileResource tsFileResource, long timeOffset, long ptNum, long valueOffset) diff --git a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTest.java b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTest.java index c573611..1862654 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTest.java @@ -310,7 +310,6 @@ public class SizeTieredCompactionRecoverTest { resourceFile.delete(); } FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); - FileReaderManager.getInstance().stop(); } /** Target file uncompleted, source files and log exists */ diff --git a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionTest.java b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionTest.java index ad4b53d..d5a50ff 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionTest.java @@ -202,7 +202,6 @@ public class SizeTieredCompactionTest { resourceFile.delete(); } FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); - FileReaderManager.getInstance().stop(); } void prepareFile(TsFileResource tsFileResource, long timeOffset, long ptNum, long valueOffset) diff --git a/server/src/test/java/org/apache/iotdb/db/engine/compaction/utils/CompactionClearUtils.java b/server/src/test/java/org/apache/iotdb/db/engine/compaction/utils/CompactionClearUtils.java index 19bfd7e..4e3af9e 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/compaction/utils/CompactionClearUtils.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/compaction/utils/CompactionClearUtils.java @@ -54,6 +54,5 @@ public class CompactionClearUtils { modsFile.delete(); } FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); - FileReaderManager.getInstance().stop(); } } diff --git a/server/src/test/java/org/apache/iotdb/db/query/control/FileReaderManagerTest.java b/server/src/test/java/org/apache/iotdb/db/query/control/FileReaderManagerTest.java index 0febbdd..20aad4d 100644 --- a/server/src/test/java/org/apache/iotdb/db/query/control/FileReaderManagerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/query/control/FileReaderManagerTest.java @@ -117,25 +117,14 @@ public class FileReaderManagerTest { t1.join(); t2.join(); + Thread.sleep(1000); + // Since we have closed the reader after reading the file, it should be false that the file is + // still contained by manager for (int i = 1; i <= MAX_FILE_SIZE; i++) { TsFileResource tsFile = new TsFileResource(SystemFileFactory.INSTANCE.getFile(filePath + i)); - Assert.assertTrue(manager.contains(tsFile, false)); + Assert.assertFalse(manager.contains(tsFile, false)); } - // the code below is not valid because the cacheFileReaderClearPeriod config in this class is - // not valid - - // TimeUnit.SECONDS.sleep(5); - // - // for (int i = 1; i <= MAX_FILE_SIZE; i++) { - // - // if (i == 4 || i == 5 || i == 6) { - // Assert.assertTrue(manager.contains(filePath + i)); - // } else { - // Assert.assertFalse(manager.contains(filePath + i)); - // } - // } - FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); for (int i = 1; i < MAX_FILE_SIZE; i++) { File file = SystemFileFactory.INSTANCE.getFile(filePath + i); diff --git a/server/src/test/java/org/apache/iotdb/db/query/reader/series/SeriesReaderTestUtil.java b/server/src/test/java/org/apache/iotdb/db/query/reader/series/SeriesReaderTestUtil.java index 9b50d67..c862482 100644 --- a/server/src/test/java/org/apache/iotdb/db/query/reader/series/SeriesReaderTestUtil.java +++ b/server/src/test/java/org/apache/iotdb/db/query/reader/series/SeriesReaderTestUtil.java @@ -205,6 +205,5 @@ public class SeriesReaderTestUtil { } FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); - FileReaderManager.getInstance().stop(); } } diff --git a/server/src/test/java/org/apache/iotdb/db/rescon/ResourceManagerTest.java b/server/src/test/java/org/apache/iotdb/db/rescon/ResourceManagerTest.java index 4c04342..27dbfe3 100644 --- a/server/src/test/java/org/apache/iotdb/db/rescon/ResourceManagerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/rescon/ResourceManagerTest.java @@ -149,7 +149,6 @@ public class ResourceManagerTest { resourceFile.delete(); } FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders(); - FileReaderManager.getInstance().stop(); } void prepareFile(TsFileResource tsFileResource, long timeOffset, long ptNum, long valueOffset)
