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();
+ }
}