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

haonan 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 03f1a52  fix cross space compaction loss data bug (#3516)
03f1a52 is described below

commit 03f1a52f5943fcdd1385033246796f2fadff5926
Author: zhanglingzhe0820 <[email protected]>
AuthorDate: Mon Jul 12 18:10:00 2021 +0800

    fix cross space compaction loss data bug (#3516)
    
    Co-authored-by: zhanglingzhe <[email protected]>
---
 .../db/engine/merge/task/MergeMultiChunkTask.java  |  11 +-
 .../iotdb/db/engine/merge/MergeTaskTest.java       | 420 +++++++++++++++++----
 2 files changed, 356 insertions(+), 75 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/task/MergeMultiChunkTask.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/task/MergeMultiChunkTask.java
index 25f17b3..d5b9811 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/task/MergeMultiChunkTask.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/task/MergeMultiChunkTask.java
@@ -206,12 +206,9 @@ public class MergeMultiChunkTask {
 
   private void pathsMergeOneFile(int seqFileIdx, IPointReader[] unseqReaders) 
throws IOException {
     TsFileResource currTsFile = resource.getSeqFiles().get(seqFileIdx);
+    // all paths in one call are from the same device
     String deviceId = currMergingPaths.get(0).getDevice();
     long currDeviceMinTime = currTsFile.getStartTime(deviceId);
-    // all paths in one call are from the same device
-    if (currDeviceMinTime == Long.MAX_VALUE) {
-      return;
-    }
 
     for (PartialPath path : currMergingPaths) {
       mergeContext.getUnmergedChunkStartTimes().get(currTsFile).put(path, new 
ArrayList<>());
@@ -335,10 +332,8 @@ public class MergeMultiChunkTask {
     // series should also be written into a new chunk
     List<Integer> ret = new ArrayList<>();
     for (int i = 0; i < currMergingPaths.size(); i++) {
-      if (seqChunkMeta[i] == null
-          || seqChunkMeta[i].isEmpty()
-              && !(seqFileIdx + 1 == resource.getSeqFiles().size()
-                  && currTimeValuePairs[i] != null)) {
+      if ((seqChunkMeta[i] == null || seqChunkMeta[i].isEmpty())
+          && !(seqFileIdx + 1 == resource.getSeqFiles().size() && 
currTimeValuePairs[i] != null)) {
         continue;
       }
       ret.add(i);
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/merge/MergeTaskTest.java 
b/server/src/test/java/org/apache/iotdb/db/engine/merge/MergeTaskTest.java
index dded161..df479d9 100644
--- a/server/src/test/java/org/apache/iotdb/db/engine/merge/MergeTaskTest.java
+++ b/server/src/test/java/org/apache/iotdb/db/engine/merge/MergeTaskTest.java
@@ -20,12 +20,17 @@
 package org.apache.iotdb.db.engine.merge;
 
 import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.fail;
 
 import java.io.File;
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.Map.Entry;
 import org.apache.commons.io.FileUtils;
+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;
@@ -33,14 +38,23 @@ import org.apache.iotdb.db.engine.merge.task.MergeTask;
 import org.apache.iotdb.db.engine.modification.Deletion;
 import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
 import org.apache.iotdb.db.exception.StorageEngineException;
+import org.apache.iotdb.db.exception.metadata.IllegalPathException;
 import org.apache.iotdb.db.exception.metadata.MetadataException;
 import org.apache.iotdb.db.metadata.PartialPath;
 import org.apache.iotdb.db.query.context.QueryContext;
 import org.apache.iotdb.db.query.reader.series.SeriesRawDataBatchReader;
 import org.apache.iotdb.tsfile.common.constant.TsFileConstant;
 import org.apache.iotdb.tsfile.exception.write.WriteProcessException;
+import org.apache.iotdb.tsfile.file.metadata.ChunkMetadata;
+import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
 import org.apache.iotdb.tsfile.read.common.BatchData;
+import org.apache.iotdb.tsfile.read.common.Path;
 import org.apache.iotdb.tsfile.read.reader.IBatchReader;
+import org.apache.iotdb.tsfile.utils.Pair;
+import org.apache.iotdb.tsfile.write.TsFileWriter;
+import org.apache.iotdb.tsfile.write.record.TSRecord;
+import org.apache.iotdb.tsfile.write.record.datapoint.DataPoint;
+import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
 import org.junit.After;
 import org.junit.Before;
 import org.junit.Test;
@@ -66,18 +80,34 @@ public class MergeTaskTest extends MergeTest {
   @Test
   public void testMerge() throws Exception {
     MergeTask mergeTask =
-        new MergeTask(new MergeResource(seqResources, unseqResources), 
tempSGDir.getPath(),
-            (k, v, l) -> {
-            }, "test", false, 1, MERGE_TEST_SG);
+        new MergeTask(
+            new MergeResource(seqResources, unseqResources),
+            tempSGDir.getPath(),
+            (k, v, l) -> {},
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
     mergeTask.call();
 
     QueryContext context = new QueryContext();
-    PartialPath path = new PartialPath(
-        deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[0].getMeasurementId());
+    PartialPath path =
+        new PartialPath(
+            deviceIds[0]
+                + TsFileConstant.PATH_SEPARATOR
+                + measurementSchemas[0].getMeasurementId());
     List<TsFileResource> list = new ArrayList<>();
     list.add(seqResources.get(0));
-    IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path, 
measurementSchemas[0].getType(),
-        context, list, new ArrayList<>(), null, null, true);
+    IBatchReader tsFilesReader =
+        new SeriesRawDataBatchReader(
+            path,
+            measurementSchemas[0].getType(),
+            context,
+            list,
+            new ArrayList<>(),
+            null,
+            null,
+            true);
     while (tsFilesReader.hasNextBatch()) {
       BatchData batchData = tsFilesReader.nextBatch();
       for (int i = 0; i < batchData.length(); i++) {
@@ -92,41 +122,59 @@ public class MergeTaskTest extends MergeTest {
     List<TsFileResource> testSeqResources = seqResources.subList(0, 3);
     List<TsFileResource> testUnseqResource = unseqResources.subList(5, 6);
     MergeTask mergeTask =
-        new MergeTask(new MergeResource(testSeqResources, testUnseqResource), 
tempSGDir.getPath(),
+        new MergeTask(
+            new MergeResource(testSeqResources, testUnseqResource),
+            tempSGDir.getPath(),
             (k, v, l) -> {
               assertEquals(499, k.get(2).getEndTime("root.mergeTest.device1"));
-            }, "test", false, 1, MERGE_TEST_SG);
+            },
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
     mergeTask.call();
   }
 
   @Test
   public void testFullMerge() throws Exception {
     MergeTask mergeTask =
-        new MergeTask(new MergeResource(seqResources, unseqResources), 
tempSGDir.getPath(),
-            (k, v, l) -> {
-            }, "test",
-            true, 10, MERGE_TEST_SG);
+        new MergeTask(
+            new MergeResource(seqResources, unseqResources),
+            tempSGDir.getPath(),
+            (k, v, l) -> {},
+            "test",
+            true,
+            10,
+            MERGE_TEST_SG);
     mergeTask.call();
 
     QueryContext context = new QueryContext();
-    PartialPath path = new PartialPath(
-        deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[9].getMeasurementId());
+    PartialPath path =
+        new PartialPath(
+            deviceIds[0]
+                + TsFileConstant.PATH_SEPARATOR
+                + measurementSchemas[9].getMeasurementId());
     List<TsFileResource> list = new ArrayList<>();
     list.add(seqResources.get(0));
-    IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path,
-        measurementSchemas[9].getType(),
-        context,
-        list, new ArrayList<>(), null, null, true);
+    IBatchReader tsFilesReader =
+        new SeriesRawDataBatchReader(
+            path,
+            measurementSchemas[9].getType(),
+            context,
+            list,
+            new ArrayList<>(),
+            null,
+            null,
+            true);
     long count = 0L;
     while (tsFilesReader.hasNextBatch()) {
       BatchData batchData = tsFilesReader.nextBatch();
       for (int t = 0; t < batchData.length(); t++) {
-        assertEquals(batchData.getTimeByIndex(t) + 20000.0, 
batchData.getDoubleByIndex(t),
-            0.001);
+        assertEquals(batchData.getTimeByIndex(t) + 20000.0, 
batchData.getDoubleByIndex(t), 0.001);
         count++;
       }
     }
-    assertEquals(100,count);
+    assertEquals(100, count);
     tsFilesReader.close();
   }
 
@@ -134,20 +182,34 @@ public class MergeTaskTest extends MergeTest {
   public void testChunkNumThreshold() throws Exception {
     
IoTDBDescriptor.getInstance().getConfig().setMergeChunkPointNumberThreshold(Integer.MAX_VALUE);
     MergeTask mergeTask =
-        new MergeTask(new MergeResource(seqResources, unseqResources), 
tempSGDir.getPath(),
-            (k, v, l) -> {
-            }, "test",
-            false, 1, MERGE_TEST_SG);
+        new MergeTask(
+            new MergeResource(seqResources, unseqResources),
+            tempSGDir.getPath(),
+            (k, v, l) -> {},
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
     mergeTask.call();
 
     QueryContext context = new QueryContext();
-    PartialPath path = new PartialPath(
-        deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[0].getMeasurementId());
+    PartialPath path =
+        new PartialPath(
+            deviceIds[0]
+                + TsFileConstant.PATH_SEPARATOR
+                + measurementSchemas[0].getMeasurementId());
     List<TsFileResource> resources = new ArrayList<>();
     resources.add(seqResources.get(0));
-    IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path, 
measurementSchemas[0].getType(),
-        context,
-        resources, new ArrayList<>(), null, null, true);
+    IBatchReader tsFilesReader =
+        new SeriesRawDataBatchReader(
+            path,
+            measurementSchemas[0].getType(),
+            context,
+            resources,
+            new ArrayList<>(),
+            null,
+            null,
+            true);
     while (tsFilesReader.hasNextBatch()) {
       BatchData batchData = tsFilesReader.nextBatch();
       for (int i = 0; i < batchData.length(); i++) {
@@ -160,19 +222,34 @@ public class MergeTaskTest extends MergeTest {
   @Test
   public void testPartialMerge1() throws Exception {
     MergeTask mergeTask =
-        new MergeTask(new MergeResource(seqResources, 
unseqResources.subList(0, 1)),
+        new MergeTask(
+            new MergeResource(seqResources, unseqResources.subList(0, 1)),
             tempSGDir.getPath(),
-            (k, v, l) -> {
-            }, "test", false, 1, MERGE_TEST_SG);
+            (k, v, l) -> {},
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
     mergeTask.call();
 
     QueryContext context = new QueryContext();
-    PartialPath path = new PartialPath(
-        deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[0].getMeasurementId());
+    PartialPath path =
+        new PartialPath(
+            deviceIds[0]
+                + TsFileConstant.PATH_SEPARATOR
+                + measurementSchemas[0].getMeasurementId());
     List<TsFileResource> list = new ArrayList<>();
     list.add(seqResources.get(0));
-    IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path, 
measurementSchemas[0].getType(),
-        context, list, new ArrayList<>(), null, null, true);
+    IBatchReader tsFilesReader =
+        new SeriesRawDataBatchReader(
+            path,
+            measurementSchemas[0].getType(),
+            context,
+            list,
+            new ArrayList<>(),
+            null,
+            null,
+            true);
     while (tsFilesReader.hasNextBatch()) {
       BatchData batchData = tsFilesReader.nextBatch();
       for (int i = 0; i < batchData.length(); i++) {
@@ -189,19 +266,34 @@ public class MergeTaskTest extends MergeTest {
   @Test
   public void testPartialMerge2() throws Exception {
     MergeTask mergeTask =
-        new MergeTask(new MergeResource(seqResources, 
unseqResources.subList(5, 6)),
+        new MergeTask(
+            new MergeResource(seqResources, unseqResources.subList(5, 6)),
             tempSGDir.getPath(),
-            (k, v, l) -> {
-            }, "test", false, 1, MERGE_TEST_SG);
+            (k, v, l) -> {},
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
     mergeTask.call();
 
     QueryContext context = new QueryContext();
-    PartialPath path = new PartialPath(
-        deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[0].getMeasurementId());
+    PartialPath path =
+        new PartialPath(
+            deviceIds[0]
+                + TsFileConstant.PATH_SEPARATOR
+                + measurementSchemas[0].getMeasurementId());
     List<TsFileResource> list = new ArrayList<>();
     list.add(seqResources.get(0));
-    IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path, 
measurementSchemas[0].getType(),
-        context, list, new ArrayList<>(), null, null, true);
+    IBatchReader tsFilesReader =
+        new SeriesRawDataBatchReader(
+            path,
+            measurementSchemas[0].getType(),
+            context,
+            list,
+            new ArrayList<>(),
+            null,
+            null,
+            true);
     while (tsFilesReader.hasNextBatch()) {
       BatchData batchData = tsFilesReader.nextBatch();
       for (int i = 0; i < batchData.length(); i++) {
@@ -214,19 +306,34 @@ public class MergeTaskTest extends MergeTest {
   @Test
   public void testPartialMerge3() throws Exception {
     MergeTask mergeTask =
-        new MergeTask(new MergeResource(seqResources, 
unseqResources.subList(0, 5)),
+        new MergeTask(
+            new MergeResource(seqResources, unseqResources.subList(0, 5)),
             tempSGDir.getPath(),
-            (k, v, l) -> {
-            }, "test", false, 1, MERGE_TEST_SG);
+            (k, v, l) -> {},
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
     mergeTask.call();
 
     QueryContext context = new QueryContext();
-    PartialPath path = new PartialPath(
-        deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[0].getMeasurementId());
+    PartialPath path =
+        new PartialPath(
+            deviceIds[0]
+                + TsFileConstant.PATH_SEPARATOR
+                + measurementSchemas[0].getMeasurementId());
     List<TsFileResource> list = new ArrayList<>();
     list.add(seqResources.get(2));
-    IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path, 
measurementSchemas[0].getType(),
-        context, list, new ArrayList<>(), null, null, true);
+    IBatchReader tsFilesReader =
+        new SeriesRawDataBatchReader(
+            path,
+            measurementSchemas[0].getType(),
+            context,
+            list,
+            new ArrayList<>(),
+            null,
+            null,
+            true);
     while (tsFilesReader.hasNextBatch()) {
       BatchData batchData = tsFilesReader.nextBatch();
       for (int i = 0; i < batchData.length(); i++) {
@@ -244,14 +351,19 @@ public class MergeTaskTest extends MergeTest {
   public void mergeWithDeletionTest() throws Exception {
     try {
       PartialPath device = new PartialPath(deviceIds[0]);
-      seqResources.get(0).getModFile().write(
-          new 
Deletion(device.concatNode(measurementSchemas[0].getMeasurementId()), 10000, 0, 
49));
+      seqResources
+          .get(0)
+          .getModFile()
+          .write(
+              new Deletion(
+                  device.concatNode(measurementSchemas[0].getMeasurementId()), 
10000, 0, 49));
     } finally {
       seqResources.get(0).getModFile().close();
     }
 
     MergeTask mergeTask =
-        new MergeTask(new MergeResource(seqResources, 
unseqResources.subList(0, 1)),
+        new MergeTask(
+            new MergeResource(seqResources, unseqResources.subList(0, 1)),
             tempSGDir.getPath(),
             (k, v, l) -> {
               try {
@@ -259,16 +371,31 @@ public class MergeTaskTest extends MergeTest {
               } catch (IOException e) {
                 e.printStackTrace();
               }
-            }, "test", false, 1, MERGE_TEST_SG);
+            },
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
     mergeTask.call();
 
     QueryContext context = new QueryContext();
-    PartialPath path = new PartialPath(
-        deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[0].getMeasurementId());
+    PartialPath path =
+        new PartialPath(
+            deviceIds[0]
+                + TsFileConstant.PATH_SEPARATOR
+                + measurementSchemas[0].getMeasurementId());
     List<TsFileResource> resources = new ArrayList<>();
     resources.add(seqResources.get(0));
-    IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path, 
measurementSchemas[0].getType(),
-        context, resources, new ArrayList<>(), null, null, true);
+    IBatchReader tsFilesReader =
+        new SeriesRawDataBatchReader(
+            path,
+            measurementSchemas[0].getType(),
+            context,
+            resources,
+            new ArrayList<>(),
+            null,
+            null,
+            true);
     int count = 0;
     while (tsFilesReader.hasNextBatch()) {
       BatchData batchData = tsFilesReader.nextBatch();
@@ -291,18 +418,34 @@ public class MergeTaskTest extends MergeTest {
     List<TsFileResource> testSeqResources = new ArrayList<>();
     List<TsFileResource> testUnseqResource = unseqResources.subList(5, 6);
     MergeTask mergeTask =
-            new MergeTask(new MergeResource(testSeqResources, 
testUnseqResource), tempSGDir.getPath(),
-                    (k, v, l) -> {
-                    }, "test", false, 1, MERGE_TEST_SG);
+        new MergeTask(
+            new MergeResource(testSeqResources, testUnseqResource),
+            tempSGDir.getPath(),
+            (k, v, l) -> {},
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
     mergeTask.call();
 
     QueryContext context = new QueryContext();
-    PartialPath path = new PartialPath(
-            deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[0].getMeasurementId());
+    PartialPath path =
+        new PartialPath(
+            deviceIds[0]
+                + TsFileConstant.PATH_SEPARATOR
+                + measurementSchemas[0].getMeasurementId());
     List<TsFileResource> resources = new ArrayList<>();
     resources.add(seqResources.get(2));
-    IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path, 
measurementSchemas[0].getType(),
-            context, resources, new ArrayList<>(), null, null, true);
+    IBatchReader tsFilesReader =
+        new SeriesRawDataBatchReader(
+            path,
+            measurementSchemas[0].getType(),
+            context,
+            resources,
+            new ArrayList<>(),
+            null,
+            null,
+            true);
     int count = 0;
     while (tsFilesReader.hasNextBatch()) {
       BatchData batchData = tsFilesReader.nextBatch();
@@ -312,4 +455,147 @@ public class MergeTaskTest extends MergeTest {
     }
     tsFilesReader.close();
   }
+
+  /**
+   * merge 3 seqFile and 1 unseqFile seqFile1: d1.s1:0-100 d1.s2:0-100 
seqFile2: d1.s1:100-200
+   * seqFile3: d2.s1:0-100 unseqFile1: d1.s3:0-100
+   */
+  @Test
+  public void testMergeWithSeqFileMissSomeSensorAndDevice() throws Exception {
+    List<TsFileResource> testSeqResources = new ArrayList<>();
+    List<TsFileResource> testUnseqResources = new ArrayList<>();
+
+    File file =
+        new File(
+            TestConstant.BASE_OUTPUT_PATH.concat(
+                100
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 100
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + ".tsfile"));
+    TsFileResource seqTsFile1 = new TsFileResource(file);
+    Map<Pair<String, MeasurementSchema>, Pair<Long, Long>> seqTsFile1Data = 
new HashMap<>();
+    seqTsFile1Data.put(new Pair<>(deviceIds[0], measurementSchemas[0]), new 
Pair<>(0L, 100L));
+    seqTsFile1Data.put(new Pair<>(deviceIds[0], measurementSchemas[1]), new 
Pair<>(0L, 100L));
+    prepareFileWithSensorAndTime(seqTsFile1, seqTsFile1Data);
+    testSeqResources.add(seqTsFile1);
+
+    file =
+        new File(
+            TestConstant.BASE_OUTPUT_PATH.concat(
+                101
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 101
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + ".tsfile"));
+    TsFileResource seqTsFile2 = new TsFileResource(file);
+    Map<Pair<String, MeasurementSchema>, Pair<Long, Long>> seqTsFileData2 = 
new HashMap<>();
+    seqTsFileData2.put(new Pair<>(deviceIds[0], measurementSchemas[0]), new 
Pair<>(100L, 200L));
+    prepareFileWithSensorAndTime(seqTsFile2, seqTsFileData2);
+    testSeqResources.add(seqTsFile2);
+
+    file =
+        new File(
+            TestConstant.BASE_OUTPUT_PATH.concat(
+                102
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 102
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + ".tsfile"));
+    TsFileResource seqTsFile3 = new TsFileResource(file);
+    Map<Pair<String, MeasurementSchema>, Pair<Long, Long>> seqTsFileData3 = 
new HashMap<>();
+    seqTsFileData3.put(new Pair<>(deviceIds[1], measurementSchemas[0]), new 
Pair<>(0L, 100L));
+    prepareFileWithSensorAndTime(seqTsFile3, seqTsFileData3);
+    testSeqResources.add(seqTsFile3);
+
+    file =
+        new File(
+            TestConstant.BASE_OUTPUT_PATH.concat(
+                10
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 10
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 10
+                    + IoTDBConstant.FILE_NAME_SEPARATOR
+                    + 0
+                    + ".tsfile"));
+    TsFileResource unseqTsFile = new TsFileResource(file);
+    Map<Pair<String, MeasurementSchema>, Pair<Long, Long>> unseqTsFileData = 
new HashMap<>();
+    unseqTsFileData.put(new Pair<>(deviceIds[0], measurementSchemas[2]), new 
Pair<>(0L, 100L));
+    prepareFileWithSensorAndTime(unseqTsFile, unseqTsFileData);
+    testUnseqResources.add(unseqTsFile);
+
+    MergeTask mergeTask =
+        new MergeTask(
+            new MergeResource(testSeqResources, testUnseqResources),
+            tempSGDir.getPath(),
+            (k, v, l) -> {
+              try (TsFileSequenceReader reader =
+                  new TsFileSequenceReader(k.get(2).getTsFilePath())) {
+                List<ChunkMetadata> chunkMetadataList =
+                    reader.getChunkMetadataList(
+                        new PartialPath(deviceIds[0], 
measurementSchemas[2].getMeasurementId()));
+                assertEquals(1, chunkMetadataList.size());
+              } catch (IOException | IllegalPathException e) {
+                e.printStackTrace();
+                fail();
+              }
+              for (TsFileResource tsFileResource : k) {
+                tsFileResource.remove();
+              }
+            },
+            "test",
+            false,
+            1,
+            MERGE_TEST_SG);
+    mergeTask.call();
+  }
+
+  /**
+   * @param tsFileResource The File to write
+   * @param generateMap map((device, measurement),(startTime, pointNum))
+   */
+  private void prepareFileWithSensorAndTime(
+      TsFileResource tsFileResource,
+      Map<Pair<String, MeasurementSchema>, Pair<Long, Long>> generateMap)
+      throws IOException, WriteProcessException {
+    TsFileWriter fileWriter = new TsFileWriter(tsFileResource.getTsFile());
+    for (Pair<String, MeasurementSchema> measurementSchemaPair : 
generateMap.keySet()) {
+      fileWriter.registerTimeseries(
+          new Path(measurementSchemaPair.left, 
measurementSchemaPair.right.getMeasurementId()),
+          measurementSchemaPair.right);
+    }
+
+    for (Entry<Pair<String, MeasurementSchema>, Pair<Long, Long>> 
generateEntry :
+        generateMap.entrySet()) {
+      Pair<String, MeasurementSchema> measurementSchemaPair = 
generateEntry.getKey();
+      String device = measurementSchemaPair.left;
+      MeasurementSchema measurementSchema = measurementSchemaPair.right;
+      Pair<Long, Long> startTimePointNumPair = generateEntry.getValue();
+      for (long i = 0; i < startTimePointNumPair.right; i++) {
+        TSRecord record = new TSRecord(i, device);
+        record.addTuple(
+            DataPoint.getDataPoint(
+                measurementSchema.getType(),
+                measurementSchema.getMeasurementId(),
+                String.valueOf(i + startTimePointNumPair.left)));
+        fileWriter.write(record);
+        tsFileResource.updateStartTime(device, i);
+        tsFileResource.updateEndTime(device, i);
+        if ((i + 1) % flushInterval == 0) {
+          fileWriter.flushAllChunkGroups();
+        }
+      }
+    }
+    fileWriter.close();
+  }
 }

Reply via email to