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