This is an automated email from the ASF dual-hosted git repository.

qiaojialin pushed a commit to branch rel/0.11
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/rel/0.11 by this push:
     new 8eeaeeb  cherry-pick [ISSUE-2746] Fix data overlapped bug after unseq 
compaction (#2750)
8eeaeeb is described below

commit 8eeaeeb08013182c5519c7d4f96656c4f5ab5407
Author: zhanglingzhe0820 <[email protected]>
AuthorDate: Mon Mar 1 17:40:39 2021 +0800

    cherry-pick [ISSUE-2746] Fix data overlapped bug after unseq compaction 
(#2750)
---
 .../db/engine/merge/manage/MergeResource.java      |   5 +-
 .../merge/selector/MaxFileMergeFileSelector.java   |   4 +-
 .../db/engine/storagegroup/TsFileResource.java     |   5 +
 .../engine/merge/MaxFileMergeFileSelectorTest.java | 127 +++++++++++++++++++--
 4 files changed, 131 insertions(+), 10 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeResource.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeResource.java
index c6baf3c..d304d4c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeResource.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeResource.java
@@ -74,8 +74,11 @@ public class MergeResource {
         .collect(Collectors.toList());
   }
 
+  /** If returns true, it means to participate in the merge */
   private boolean filterResource(TsFileResource res) {
-    return res.getTsFile().exists() && !res.isDeleted() && 
res.stillLives(timeLowerBound);
+    return res.getTsFile().exists()
+        && !res.isDeleted()
+        && (!res.isClosed() || res.stillLives(timeLowerBound));
   }
 
   public MergeResource(Collection<TsFileResource> seqFiles, 
List<TsFileResource> unseqFiles,
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/selector/MaxFileMergeFileSelector.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/selector/MaxFileMergeFileSelector.java
index d115a11..7cc23e5 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/selector/MaxFileMergeFileSelector.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/selector/MaxFileMergeFileSelector.java
@@ -211,7 +211,6 @@ public class MaxFileMergeFileSelector implements 
IMergeFileSelector {
   }
 
   private void selectOverlappedSeqFiles(TsFileResource unseqFile) {
-
     int tmpSelectedNum = 0;
     for (Entry<String, Integer> deviceStartTimeEntry : 
unseqFile.getDeviceToIndexMap().entrySet()) {
       String deviceId = deviceStartTimeEntry.getKey();
@@ -225,7 +224,8 @@ public class MaxFileMergeFileSelector implements 
IMergeFileSelector {
         if (seqSelected[i] || 
!seqFile.getDeviceToIndexMap().containsKey(deviceId)) {
           continue;
         }
-        long seqEndTime = seqFile.getEndTime(deviceId);
+        // the open file's endTime is Long.MIN_VALUE, this will make the file 
be filtered below
+        long seqEndTime = seqFile.isClosed() ? seqFile.getEndTime(deviceId) : 
Long.MAX_VALUE;
         if (unseqEndTime <= seqEndTime) {
           // the unseqFile overlaps current seqFile
           tmpSelectedSeqFiles.add(i);
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 f136732..559bd01 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
@@ -400,6 +400,7 @@ public class TsFileResource {
     return startTimes[index];
   }
 
+  /** open file's end time is Long.MIN_VALUE */
   public long getEndTime(String deviceId) {
     if (!deviceToIndex.containsKey(deviceId)) {
       return Long.MIN_VALUE;
@@ -892,4 +893,8 @@ public class TsFileResource {
     return another.maxPlanIndex >= this.minPlanIndex &&
            another.minPlanIndex <= this.maxPlanIndex;
   }
+
+  public void setTimeIndex(ITimeIndex timeIndex) {
+    this.timeIndex = timeIndex;
+  }
 }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxFileMergeFileSelectorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxFileMergeFileSelectorTest.java
index 1e50590..7275a01 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxFileMergeFileSelectorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/merge/MaxFileMergeFileSelectorTest.java
@@ -19,17 +19,27 @@
 
 package org.apache.iotdb.db.engine.merge;
 
-import static org.junit.Assert.assertEquals;
-
-import java.io.IOException;
-import java.util.List;
+import org.apache.iotdb.db.conf.IoTDBConstant;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.constant.TestConstant;
 import org.apache.iotdb.db.engine.merge.manage.MergeResource;
-import org.apache.iotdb.db.engine.merge.selector.MaxFileMergeFileSelector;
 import org.apache.iotdb.db.engine.merge.selector.IMergeFileSelector;
+import org.apache.iotdb.db.engine.merge.selector.MaxFileMergeFileSelector;
 import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.db.engine.storagegroup.timeindex.ITimeIndex;
 import org.apache.iotdb.db.exception.MergeException;
+import org.apache.iotdb.tsfile.exception.write.WriteProcessException;
+
 import org.junit.Test;
 
+import java.io.File;
+import java.io.IOException;
+import java.lang.reflect.Field;
+import java.util.List;
+import java.util.Set;
+
+import static org.junit.Assert.assertEquals;
+
 public class MaxFileMergeFileSelectorTest extends MergeTest {
 
   @Test
@@ -78,8 +88,111 @@ public class MaxFileMergeFileSelectorTest extends MergeTest 
{
     List[] result = mergeFileSelector.select();
     List<TsFileResource> seqSelected = result[0];
     List<TsFileResource> unseqSelected = result[1];
-    assertEquals(seqResources.subList(0, 3), seqSelected);
-    assertEquals(unseqResources.subList(0, 3), unseqSelected);
+    assertEquals(seqResources.subList(0, 4), seqSelected);
+    assertEquals(unseqResources.subList(0, 4), unseqSelected);
     resource.clear();
   }
+
+  /**
+   * test unseq merge select with the following files: {0seq-0-0-0.tsfile 
0-100 1seq-1-1-0.tsfile
+   * 100-200 2seq-2-2-0.tsfile 200-300 3seq-3-3-0.tsfile 300-400 
4seq-4-4-0.tsfile 400-500}
+   * {10unseq-10-10-0.tsfile 0-500}
+   */
+  @Test
+  public void testFileOpenSelection()
+      throws MergeException, IOException, WriteProcessException, 
NoSuchFieldException,
+          IllegalAccessException {
+    File file =
+        new File(
+            TestConstant.BASE_OUTPUT_PATH.concat(
+                10
+                    + "unseq"
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 10
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 10
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + ".tsfile"));
+    TsFileResource largeUnseqTsFileResource = new TsFileResource(file);
+    largeUnseqTsFileResource.setClosed(true);
+    largeUnseqTsFileResource.setMinPlanIndex(10);
+    largeUnseqTsFileResource.setMaxPlanIndex(10);
+    largeUnseqTsFileResource.setVersion(10);
+    prepareFile(largeUnseqTsFileResource, 0, seqFileNum * ptNum, 0);
+
+    // update the second file's status to open
+    TsFileResource secondTsFileResource = seqResources.get(1);
+    secondTsFileResource.setClosed(false);
+    Set<String> devices = secondTsFileResource.getDevices();
+    // update the end time of the file to Long.MIN_VALUE, so we can simulate a 
real open file
+    Field timeIndexField = TsFileResource.class.getDeclaredField("timeIndex");
+    timeIndexField.setAccessible(true);
+    ITimeIndex timeIndex = (ITimeIndex) 
timeIndexField.get(secondTsFileResource);
+    ITimeIndex newTimeIndex =
+        
IoTDBDescriptor.getInstance().getConfig().getTimeIndexLevel().getTimeIndex();
+    for (String device : devices) {
+      newTimeIndex.updateStartTime(device, timeIndex.getStartTime(device));
+    }
+    secondTsFileResource.setTimeIndex(newTimeIndex);
+    unseqResources.clear();
+    unseqResources.add(largeUnseqTsFileResource);
+
+    MergeResource resource = new MergeResource(seqResources, unseqResources);
+    IMergeFileSelector mergeFileSelector = new 
MaxFileMergeFileSelector(resource, Long.MAX_VALUE);
+    List[] result = mergeFileSelector.select();
+    assertEquals(0, result.length);
+    resource.clear();
+  }
+
+  /**
+   * test unseq merge select with the following files: {0seq-0-0-0.tsfile 
0-100 1seq-1-1-0.tsfile
+   * 100-200 2seq-2-2-0.tsfile 200-300 3seq-3-3-0.tsfile 300-400 
4seq-4-4-0.tsfile 400-500}
+   * {10unseq-10-10-0.tsfile 0-500}
+   */
+  @Test
+  public void testFileOpenSelectionFromCompaction()
+      throws IOException, WriteProcessException, NoSuchFieldException, 
IllegalAccessException {
+    File file =
+        new File(
+            TestConstant.BASE_OUTPUT_PATH.concat(
+                10
+                    + "unseq"
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 10
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 10
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + ".tsfile"));
+    TsFileResource largeUnseqTsFileResource = new TsFileResource(file);
+    largeUnseqTsFileResource.setClosed(true);
+    largeUnseqTsFileResource.setMinPlanIndex(10);
+    largeUnseqTsFileResource.setMaxPlanIndex(10);
+    largeUnseqTsFileResource.setVersion(10);
+    prepareFile(largeUnseqTsFileResource, 0, seqFileNum * ptNum, 0);
+
+    // update the second file's status to open
+    TsFileResource secondTsFileResource = seqResources.get(1);
+    secondTsFileResource.setClosed(false);
+    Set<String> devices = secondTsFileResource.getDevices();
+    // update the end time of the file to Long.MIN_VALUE, so we can simulate a 
real open file
+    Field timeIndexField = TsFileResource.class.getDeclaredField("timeIndex");
+    timeIndexField.setAccessible(true);
+    ITimeIndex timeIndex = (ITimeIndex) 
timeIndexField.get(secondTsFileResource);
+    ITimeIndex newTimeIndex =
+        
IoTDBDescriptor.getInstance().getConfig().getTimeIndexLevel().getTimeIndex();
+    for (String device : devices) {
+      newTimeIndex.updateStartTime(device, timeIndex.getStartTime(device));
+    }
+    secondTsFileResource.setTimeIndex(newTimeIndex);
+    unseqResources.clear();
+    unseqResources.add(largeUnseqTsFileResource);
+
+    long timeLowerBound = System.currentTimeMillis() - Long.MAX_VALUE;
+    MergeResource mergeResource = new MergeResource(seqResources, 
unseqResources, timeLowerBound);
+    assertEquals(5, mergeResource.getSeqFiles().size());
+    assertEquals(1, mergeResource.getUnseqFiles().size());
+    mergeResource.clear();
+  }
 }

Reply via email to