This is an automated email from the ASF dual-hosted git repository. sunzesong pushed a commit to branch time_index_improve in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit f5c999585ea12100d3e0f3378787302046114ed8 Author: samperson1997 <[email protected]> AuthorDate: Wed Dec 30 10:40:07 2020 +0800 Record the device number of the last TsFile in each storage group --- .../engine/storagegroup/StorageGroupProcessor.java | 13 +++++++++-- .../db/engine/storagegroup/TsFileProcessor.java | 7 +++--- .../db/engine/storagegroup/TsFileResource.java | 4 ++-- .../storagegroup/timeindex/DeviceTimeIndex.java | 27 ++++++++++++---------- .../storagegroup/timeindex/FileTimeIndex.java | 11 +++------ .../engine/storagegroup/timeindex/ITimeIndex.java | 5 ---- .../storagegroup/timeindex/TimeIndexLevel.java | 10 ++++++++ .../engine/storagegroup/TsFileProcessorTest.java | 10 ++++---- 8 files changed, 51 insertions(+), 36 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java index 79a4403..438e6f3 100755 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java @@ -59,6 +59,7 @@ import org.apache.iotdb.db.engine.merge.task.RecoverMergeTask; import org.apache.iotdb.db.engine.modification.Deletion; import org.apache.iotdb.db.engine.modification.ModificationFile; import org.apache.iotdb.db.engine.querycontext.QueryDataSource; +import org.apache.iotdb.db.engine.storagegroup.timeindex.DeviceTimeIndex; import org.apache.iotdb.db.engine.version.SimpleFileVersionController; import org.apache.iotdb.db.engine.version.VersionController; import org.apache.iotdb.db.exception.BatchProcessException; @@ -239,6 +240,13 @@ public class StorageGroupProcessor { private long monitorSeriesValue; private StorageGroupInfo storageGroupInfo = new StorageGroupInfo(this); + /** + * Record the device number of the last TsFile in each storage group, which is applied to + * initialize the array size of DeviceTimeIndex. It is reasonable to assume that the adjacent + * files should have similar numbers of devices. Default value: INIT_ARRAY_SIZE = 64 + */ + private int deviceNumInLastClosedTsFile = DeviceTimeIndex.INIT_ARRAY_SIZE; + public boolean isReady() { return isReady; } @@ -1020,7 +1028,7 @@ public class StorageGroupProcessor { tsFileProcessor = new TsFileProcessor(storageGroupName, fsFactory.getFileWithParent(filePath), storageGroupInfo, versionController, this::closeUnsealedTsFileProcessorCallBack, - this::updateLatestFlushTimeCallback, true); + this::updateLatestFlushTimeCallback, true, deviceNumInLastClosedTsFile); if (enableMemControl) { TsFileProcessorInfo tsFileProcessorInfo = new TsFileProcessorInfo(storageGroupInfo); tsFileProcessor.setTsFileProcessorInfo(tsFileProcessorInfo); @@ -1032,7 +1040,7 @@ public class StorageGroupProcessor { tsFileProcessor = new TsFileProcessor(storageGroupName, fsFactory.getFileWithParent(filePath), storageGroupInfo, versionController, this::closeUnsealedTsFileProcessorCallBack, - this::unsequenceFlushCallback, false); + this::unsequenceFlushCallback, false, deviceNumInLastClosedTsFile); if (enableMemControl) { TsFileProcessorInfo tsFileProcessorInfo = new TsFileProcessorInfo(storageGroupInfo); tsFileProcessor.setTsFileProcessorInfo(tsFileProcessorInfo); @@ -1646,6 +1654,7 @@ public class StorageGroupProcessor { closeQueryLock.writeLock().lock(); try { tsFileProcessor.close(); + deviceNumInLastClosedTsFile = tsFileProcessor.getTsFileResource().getDevices().size(); } finally { closeQueryLock.writeLock().unlock(); } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java index 18c16d9..9ba9a75 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java @@ -127,10 +127,11 @@ public class TsFileProcessor { StorageGroupInfo storageGroupInfo, VersionController versionController, CloseFileListener closeTsFileCallback, - UpdateEndTimeCallBack updateLatestFlushTimeCallback, boolean sequence) + UpdateEndTimeCallBack updateLatestFlushTimeCallback, boolean sequence, + int deviceNumInLastClosedTsFile) throws IOException { this.storageGroupName = storageGroupName; - this.tsFileResource = new TsFileResource(tsfile, this); + this.tsFileResource = new TsFileResource(tsfile, this, deviceNumInLastClosedTsFile); if (enableMemControl) { this.storageGroupInfo = storageGroupInfo; } @@ -861,7 +862,7 @@ public class TsFileProcessor { public void close() throws TsFileProcessorException { try { - //when closing resource file, its corresponding mod file is also closed. + // when closing resource file, its corresponding mod file is also closed. tsFileResource.close(); MultiFileLogNodeManager.getInstance() .deleteNode(storageGroupName + "-" + tsFileResource.getTsFile().getName()); diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java index 07ec55d..a9fa5bc 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java @@ -171,9 +171,9 @@ public class TsFileResource { /** * unsealed TsFile */ - public TsFileResource(File file, TsFileProcessor processor) { + public TsFileResource(File file, TsFileProcessor processor, int deviceNumInLastClosedTsFile) { this.file = file; - this.timeIndex = config.getTimeIndexLevel().getTimeIndex(); + this.timeIndex = config.getTimeIndexLevel().getTimeIndex(deviceNumInLastClosedTsFile); this.processor = processor; } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/DeviceTimeIndex.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/DeviceTimeIndex.java index c3ebb4f..13f8313 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/DeviceTimeIndex.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/DeviceTimeIndex.java @@ -38,7 +38,7 @@ import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils; public class DeviceTimeIndex implements ITimeIndex { - protected static final int INIT_ARRAY_SIZE = 64; + public static final int INIT_ARRAY_SIZE = 64; protected static final Map<String, String> cachedDevicePool = CachedStringPool.getInstance() .getCachedPool(); @@ -60,7 +60,19 @@ public class DeviceTimeIndex implements ITimeIndex { protected Map<String, Integer> deviceToIndex; public DeviceTimeIndex() { - init(); + this.deviceToIndex = new ConcurrentHashMap<>(); + this.startTimes = new long[INIT_ARRAY_SIZE]; + this.endTimes = new long[INIT_ARRAY_SIZE]; + initTimes(startTimes, Long.MAX_VALUE); + initTimes(endTimes, Long.MIN_VALUE); + } + + public DeviceTimeIndex(int deviceNumInLastClosedTsFile) { + this.deviceToIndex = new ConcurrentHashMap<>(); + this.startTimes = new long[deviceNumInLastClosedTsFile]; + this.endTimes = new long[deviceNumInLastClosedTsFile]; + initTimes(startTimes, Long.MAX_VALUE); + initTimes(endTimes, Long.MIN_VALUE); } public DeviceTimeIndex(Map<String, Integer> deviceToIndex, long[] startTimes, long[] endTimes) { @@ -70,15 +82,6 @@ public class DeviceTimeIndex implements ITimeIndex { } @Override - public void init() { - this.deviceToIndex = new ConcurrentHashMap<>(); - this.startTimes = new long[INIT_ARRAY_SIZE]; - this.endTimes = new long[INIT_ARRAY_SIZE]; - initTimes(startTimes, Long.MAX_VALUE); - initTimes(endTimes, Long.MIN_VALUE); - } - - @Override public void serialize(OutputStream outputStream) throws IOException { int deviceNum = deviceToIndex.size(); @@ -221,7 +224,7 @@ public class DeviceTimeIndex implements ITimeIndex { } private long[] enLargeArray(long[] array, long defaultValue) { - long[] tmp = new long[(int) (array.length * 1.5)]; + long[] tmp = new long[(int) (array.length * 2)]; initTimes(tmp, defaultValue); System.arraycopy(array, 0, tmp, 0, array.length); return tmp; diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/FileTimeIndex.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/FileTimeIndex.java index b51ec42..c8d68f7 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/FileTimeIndex.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/FileTimeIndex.java @@ -56,7 +56,9 @@ public class FileTimeIndex implements ITimeIndex { protected Set<String> devices; public FileTimeIndex() { - init(); + this.devices = new ConcurrentSet<>(); + this.startTime = Long.MAX_VALUE; + this.endTime = Long.MIN_VALUE; } public FileTimeIndex(Set<String> devices, long startTime, long endTime) { @@ -66,13 +68,6 @@ public class FileTimeIndex implements ITimeIndex { } @Override - public void init() { - this.devices = new ConcurrentSet<>(); - this.startTime = Long.MAX_VALUE; - this.endTime = Long.MIN_VALUE; - } - - @Override public void serialize(OutputStream outputStream) throws IOException { ReadWriteIOUtils.write(devices.size(), outputStream); for (String device : devices) { diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/ITimeIndex.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/ITimeIndex.java index afeb91d..fda0d32 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/ITimeIndex.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/ITimeIndex.java @@ -29,11 +29,6 @@ import org.apache.iotdb.db.exception.PartitionViolationException; public interface ITimeIndex { /** - * init startTimes with Long.MAX_VALUE, endTimes with Long.MIN_VALUE - */ - void init(); - - /** * serialize to outputStream * * @param outputStream outputStream diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/TimeIndexLevel.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/TimeIndexLevel.java index e1a80b0..7f3b624 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/TimeIndexLevel.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/timeindex/TimeIndexLevel.java @@ -31,4 +31,14 @@ public enum TimeIndexLevel { return new DeviceTimeIndex(); } } + + public ITimeIndex getTimeIndex(int deviceNumInLastClosedTsFile) { + switch (this) { + case FILE_TIME_INDEX: + return new FileTimeIndex(); + case DEVICE_TIME_INDEX: + default: + return new DeviceTimeIndex(deviceNumInLastClosedTsFile); + } + } } diff --git a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessorTest.java b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessorTest.java index 362ba9c..d819780 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessorTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessorTest.java @@ -69,6 +69,8 @@ public class TsFileProcessorTest { private QueryContext context; private static Logger logger = LoggerFactory.getLogger(TsFileProcessorTest.class); + protected static final int INIT_ARRAY_SIZE = 64; + @Before public void setUp() throws Exception { EnvironmentUtils.envSetUp(); @@ -87,7 +89,7 @@ public class TsFileProcessorTest { logger.info("testWriteAndFlush begin.."); processor = new TsFileProcessor(storageGroup, SystemFileFactory.INSTANCE.getFile(filePath), sgInfo, SysTimeVersionController.INSTANCE, this::closeTsFileProcessor, - (tsFileProcessor) -> true, true); + (tsFileProcessor) -> true, true, INIT_ARRAY_SIZE); TsFileProcessorInfo tsFileProcessorInfo = new TsFileProcessorInfo(sgInfo); processor.setTsFileProcessorInfo(tsFileProcessorInfo); @@ -143,7 +145,7 @@ public class TsFileProcessorTest { logger.info("testWriteAndRestoreMetadata begin.."); processor = new TsFileProcessor(storageGroup, SystemFileFactory.INSTANCE.getFile(filePath), sgInfo, SysTimeVersionController.INSTANCE, this::closeTsFileProcessor, - (tsFileProcessor) -> true, true); + (tsFileProcessor) -> true, true, INIT_ARRAY_SIZE); TsFileProcessorInfo tsFileProcessorInfo = new TsFileProcessorInfo(sgInfo); processor.setTsFileProcessorInfo(tsFileProcessorInfo); @@ -226,7 +228,7 @@ public class TsFileProcessorTest { logger.info("testWriteAndRestoreMetadata begin.."); processor = new TsFileProcessor(storageGroup, SystemFileFactory.INSTANCE.getFile(filePath), sgInfo, SysTimeVersionController.INSTANCE, this::closeTsFileProcessor, - (tsFileProcessor) -> true, true); + (tsFileProcessor) -> true, true, INIT_ARRAY_SIZE); TsFileProcessorInfo tsFileProcessorInfo = new TsFileProcessorInfo(sgInfo); processor.setTsFileProcessorInfo(tsFileProcessorInfo); @@ -267,7 +269,7 @@ public class TsFileProcessorTest { logger.info("testWriteAndRestoreMetadata begin.."); processor = new TsFileProcessor(storageGroup, SystemFileFactory.INSTANCE.getFile(filePath), sgInfo, SysTimeVersionController.INSTANCE, this::closeTsFileProcessor, - (tsFileProcessor) -> true, true); + (tsFileProcessor) -> true, true, INIT_ARRAY_SIZE); TsFileProcessorInfo tsFileProcessorInfo = new TsFileProcessorInfo(sgInfo); processor.setTsFileProcessorInfo(tsFileProcessorInfo);
