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

Reply via email to