This is an automated email from the ASF dual-hosted git repository.
haonan pushed a commit to branch rel/0.13
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rel/0.13 by this push:
new c0d301516a [To rel/0.13][IOTDB-3247] Recover aligned sensors after
deleting timeseries, query lost data (#6548)
c0d301516a is described below
commit c0d301516adb164b9f5f6e3b159baa3b1bfa172a
Author: Chen YZ <[email protected]>
AuthorDate: Tue Jul 12 15:31:49 2022 +0800
[To rel/0.13][IOTDB-3247] Recover aligned sensors after deleting
timeseries, query lost data (#6548)
---
.../iotdb/db/engine/memtable/AbstractMemTable.java | 22 ++-
.../engine/memtable/AlignedWritableMemChunk.java | 46 ++++--
.../memtable/AlignedWritableMemChunkGroup.java | 11 +-
.../apache/iotdb/db/engine/memtable/IMemTable.java | 2 +
.../db/engine/memtable/IWritableMemChunk.java | 6 +-
.../db/engine/memtable/IWritableMemChunkGroup.java | 7 +-
.../iotdb/db/engine/memtable/WritableMemChunk.java | 6 +-
.../db/engine/memtable/WritableMemChunkGroup.java | 7 +-
.../iotdb/db/qp/physical/crud/InsertPlan.java | 5 +
.../db/engine/memtable/MemTableTestUtils.java | 1 +
.../db/engine/memtable/MemtableBenchmark.java | 1 +
.../db/engine/memtable/PrimitiveMemTableTest.java | 11 ++
.../iotdb/db/writelog/recover/LogReplayerTest.java | 156 +++++++++++++++++++--
13 files changed, 247 insertions(+), 34 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
index 7cde23546b..ab8d6b6f4d 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
@@ -168,7 +168,12 @@ public abstract class AbstractMemTable implements
IMemTable {
dataTypes.add(schema.getType());
}
memSize += MemUtils.getRecordsSize(dataTypes, values, disableMemControl);
- write(insertRowPlan.getDeviceID(), schemaList, insertRowPlan.getTime(),
values);
+ write(
+ insertRowPlan.getDeviceID(),
+ insertRowPlan.getFailedIndices(),
+ schemaList,
+ insertRowPlan.getTime(),
+ values);
int pointsInserted =
insertRowPlan.getMeasurements().length
@@ -216,7 +221,12 @@ public abstract class AbstractMemTable implements
IMemTable {
return;
}
memSize += MemUtils.getAlignedRecordsSize(dataTypes, values,
disableMemControl);
- writeAlignedRow(insertRowPlan.getDeviceID(), schemaList,
insertRowPlan.getTime(), values);
+ writeAlignedRow(
+ insertRowPlan.getDeviceID(),
+ insertRowPlan.getFailedIndices(),
+ schemaList,
+ insertRowPlan.getTime(),
+ values);
int pointsInserted =
insertRowPlan.getMeasurements().length -
insertRowPlan.getFailedMeasurementNumber();
totalPointsNum += pointsInserted;
@@ -288,23 +298,25 @@ public abstract class AbstractMemTable implements
IMemTable {
@Override
public void write(
IDeviceID deviceId,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
long insertTime,
Object[] objectValue) {
IWritableMemChunkGroup memChunkGroup =
createMemChunkGroupIfNotExistAndGet(deviceId, schemaList);
- memChunkGroup.write(insertTime, objectValue, schemaList);
+ memChunkGroup.write(insertTime, objectValue, failedIndices, schemaList);
}
@Override
public void writeAlignedRow(
IDeviceID deviceId,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
long insertTime,
Object[] objectValue) {
IWritableMemChunkGroup memChunkGroup =
createAlignedMemChunkGroupIfNotExistAndGet(deviceId, schemaList);
- memChunkGroup.write(insertTime, objectValue, schemaList);
+ memChunkGroup.write(insertTime, objectValue, failedIndices, schemaList);
}
@SuppressWarnings("squid:S3776") // high Cognitive Complexity
@@ -330,6 +342,7 @@ public abstract class AbstractMemTable implements IMemTable
{
insertTabletPlan.getTimes(),
insertTabletPlan.getColumns(),
insertTabletPlan.getBitMaps(),
+ insertTabletPlan.getFailedIndices(),
schemaList,
start,
end);
@@ -360,6 +373,7 @@ public abstract class AbstractMemTable implements IMemTable
{
insertTabletPlan.getTimes(),
insertTabletPlan.getColumns(),
insertTabletPlan.getBitMaps(),
+ insertTabletPlan.getFailedIndices(),
schemaList,
start,
end);
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AlignedWritableMemChunk.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AlignedWritableMemChunk.java
index c7ec62f1f1..eeec5511ef 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AlignedWritableMemChunk.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AlignedWritableMemChunk.java
@@ -144,8 +144,11 @@ public class AlignedWritableMemChunk implements
IWritableMemChunk {
@Override
public void writeAlignedValue(
- long insertTime, Object[] objectValue, List<IMeasurementSchema>
schemaList) {
- int[] columnIndexArray = checkColumnsInInsertPlan(schemaList);
+ long insertTime,
+ Object[] objectValue,
+ List<Integer> failedIndices,
+ List<IMeasurementSchema> schemaList) {
+ int[] columnIndexArray = checkColumnsInInsertPlan(failedIndices,
schemaList);
putAlignedValue(insertTime, objectValue, columnIndexArray);
}
@@ -160,29 +163,44 @@ public class AlignedWritableMemChunk implements
IWritableMemChunk {
long[] times,
Object[] valueList,
BitMap[] bitMaps,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
int start,
int end) {
- int[] columnIndexArray = checkColumnsInInsertPlan(schemaList);
+ int[] columnIndexArray = checkColumnsInInsertPlan(failedIndices,
schemaList);
putAlignedValues(times, valueList, bitMaps, columnIndexArray, start, end);
}
- private int[] checkColumnsInInsertPlan(List<IMeasurementSchema>
schemaListInInsertPlan) {
+ /**
+ * Check schema of columns and return array that mapping existed schema to
index of data column
+ *
+ * @param failedIndices It records the index of timeseries that have been
deleted.
+ * @param schemaListInInsertPlan Contains all schema in InsertPlan.
+ * @return columnIndexArray: schemaList[i] is schema of
columns[columnIndexArray[i]]
+ */
+ private int[] checkColumnsInInsertPlan(
+ List<Integer> failedIndices, List<IMeasurementSchema>
schemaListInInsertPlan) {
Map<String, Integer> measurementIdsInInsertPlan = new HashMap<>();
- for (int i = 0; i < schemaListInInsertPlan.size(); i++) {
-
measurementIdsInInsertPlan.put(schemaListInInsertPlan.get(i).getMeasurementId(),
i);
- if
(!containsMeasurement(schemaListInInsertPlan.get(i).getMeasurementId())) {
- this.measurementIndexMap.put(
- schemaListInInsertPlan.get(i).getMeasurementId(),
measurementIndexMap.size());
- this.schemaList.add(schemaListInInsertPlan.get(i));
- this.list.extendColumn(schemaListInInsertPlan.get(i).getType());
+ for (int i = 0, failedIndicesIdx = 0, schemaListIdx = 0;
+ i < failedIndices.size() + schemaListInInsertPlan.size();
+ i++) {
+ if (failedIndices.size() > failedIndicesIdx &&
failedIndices.get(failedIndicesIdx) == i) {
+ failedIndicesIdx++;
+ } else {
+ IMeasurementSchema measurementSchema =
schemaListInInsertPlan.get(schemaListIdx++);
+ measurementIdsInInsertPlan.put(measurementSchema.getMeasurementId(),
i);
+ if (!containsMeasurement(measurementSchema.getMeasurementId())) {
+ this.measurementIndexMap.put(
+ measurementSchema.getMeasurementId(),
measurementIndexMap.size());
+ this.schemaList.add(measurementSchema);
+ this.list.extendColumn(measurementSchema.getType());
+ }
}
}
int[] columnIndexArray = new int[measurementIndexMap.size()];
measurementIndexMap.forEach(
- (measurementId, i) -> {
- columnIndexArray[i] =
measurementIdsInInsertPlan.getOrDefault(measurementId, -1);
- });
+ (measurementId, i) ->
+ columnIndexArray[i] =
measurementIdsInInsertPlan.getOrDefault(measurementId, -1));
return columnIndexArray;
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AlignedWritableMemChunkGroup.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AlignedWritableMemChunkGroup.java
index 56fc7e9501..cfe957b86c 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AlignedWritableMemChunkGroup.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AlignedWritableMemChunkGroup.java
@@ -44,10 +44,11 @@ public class AlignedWritableMemChunkGroup implements
IWritableMemChunkGroup {
long[] times,
Object[] columns,
BitMap[] bitMaps,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
int start,
int end) {
- memChunk.writeAlignedValues(times, columns, bitMaps, schemaList, start,
end);
+ memChunk.writeAlignedValues(times, columns, bitMaps, failedIndices,
schemaList, start, end);
}
@Override
@@ -74,8 +75,12 @@ public class AlignedWritableMemChunkGroup implements
IWritableMemChunkGroup {
}
@Override
- public void write(long insertTime, Object[] objectValue,
List<IMeasurementSchema> schemaList) {
- memChunk.writeAlignedValue(insertTime, objectValue, schemaList);
+ public void write(
+ long insertTime,
+ Object[] objectValue,
+ List<Integer> failedIndices,
+ List<IMeasurementSchema> schemaList) {
+ memChunk.writeAlignedValue(insertTime, objectValue, failedIndices,
schemaList);
}
@Override
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
index 69aeb47ffe..23f19b1c5c 100644
--- a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
+++ b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
@@ -48,12 +48,14 @@ public interface IMemTable {
void write(
IDeviceID deviceId,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
long insertTime,
Object[] objectValue);
void writeAlignedRow(
IDeviceID deviceId,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
long insertTime,
Object[] objectValue);
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IWritableMemChunk.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IWritableMemChunk.java
index c46714fbba..7afed34eb8 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IWritableMemChunk.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IWritableMemChunk.java
@@ -61,7 +61,10 @@ public interface IWritableMemChunk {
void write(long insertTime, Object objectValue);
void writeAlignedValue(
- long insertTime, Object[] objectValue, List<IMeasurementSchema>
schemaList);
+ long insertTime,
+ Object[] objectValue,
+ List<Integer> failedIndices,
+ List<IMeasurementSchema> schemaList);
/**
* write data in the range [start, end). Null value in the valueList will be
replaced by the
@@ -74,6 +77,7 @@ public interface IWritableMemChunk {
long[] times,
Object[] valueList,
BitMap[] bitMaps,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
int start,
int end);
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IWritableMemChunkGroup.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IWritableMemChunkGroup.java
index 0ee70755b5..00138ff6a8 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IWritableMemChunkGroup.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IWritableMemChunkGroup.java
@@ -32,6 +32,7 @@ public interface IWritableMemChunkGroup {
long[] times,
Object[] columns,
BitMap[] bitMaps,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
int start,
int end);
@@ -42,7 +43,11 @@ public interface IWritableMemChunkGroup {
boolean contains(String measurement);
- void write(long insertTime, Object[] objectValue, List<IMeasurementSchema>
schemaList);
+ void write(
+ long insertTime,
+ Object[] objectValue,
+ List<Integer> failedIndices,
+ List<IMeasurementSchema> schemaList);
Map<String, IWritableMemChunk> getMemChunkMap();
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/WritableMemChunk.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/WritableMemChunk.java
index 5f99bb6d1a..0d961765be 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/WritableMemChunk.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/WritableMemChunk.java
@@ -72,7 +72,10 @@ public class WritableMemChunk implements IWritableMemChunk {
@Override
public void writeAlignedValue(
- long insertTime, Object[] objectValue, List<IMeasurementSchema>
schemaList) {
+ long insertTime,
+ Object[] objectValue,
+ List<Integer> failedIndices,
+ List<IMeasurementSchema> schemaList) {
throw new UnSupportedDataTypeException(UNSUPPORTED_TYPE +
list.getDataType());
}
@@ -114,6 +117,7 @@ public class WritableMemChunk implements IWritableMemChunk {
long[] times,
Object[] valueList,
BitMap[] bitMaps,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
int start,
int end) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/WritableMemChunkGroup.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/WritableMemChunkGroup.java
index 820fa41459..aae462df7a 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/WritableMemChunkGroup.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/WritableMemChunkGroup.java
@@ -42,6 +42,7 @@ public class WritableMemChunkGroup implements
IWritableMemChunkGroup {
long[] times,
Object[] columns,
BitMap[] bitMaps,
+ List<Integer> failedIndices,
List<IMeasurementSchema> schemaList,
int start,
int end) {
@@ -90,7 +91,11 @@ public class WritableMemChunkGroup implements
IWritableMemChunkGroup {
}
@Override
- public void write(long insertTime, Object[] objectValue,
List<IMeasurementSchema> schemaList) {
+ public void write(
+ long insertTime,
+ Object[] objectValue,
+ List<Integer> failedIndices,
+ List<IMeasurementSchema> schemaList) {
int emptyColumnCount = 0;
for (int i = 0; i < objectValue.length; i++) {
if (objectValue[i] == null) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertPlan.java
index 21ef5b38bf..d1943ab576 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertPlan.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertPlan.java
@@ -29,6 +29,7 @@ import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
@@ -113,6 +114,10 @@ public abstract class InsertPlan extends PhysicalPlan {
return failedMeasurements == null ? 0 : failedMeasurements.size();
}
+ public List<Integer> getFailedIndices() {
+ return failedIndices == null ? Collections.emptyList() : failedIndices;
+ }
+
public boolean isAligned() {
return isAligned;
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/memtable/MemTableTestUtils.java
b/server/src/test/java/org/apache/iotdb/db/engine/memtable/MemTableTestUtils.java
index 9ef4316a01..3a303a6a18 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/memtable/MemTableTestUtils.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/memtable/MemTableTestUtils.java
@@ -64,6 +64,7 @@ public class MemTableTestUtils {
for (long l = startTime; l <= endTime; l++) {
iMemTable.write(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId, dataType,
TSEncoding.PLAIN)),
l,
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/memtable/MemtableBenchmark.java
b/server/src/test/java/org/apache/iotdb/db/engine/memtable/MemtableBenchmark.java
index 0f60ee8b9a..96b4da1a67 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/memtable/MemtableBenchmark.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/memtable/MemtableBenchmark.java
@@ -51,6 +51,7 @@ public class MemtableBenchmark {
for (int j = 0; j < numOfMeasurement; j++) {
memTable.write(
DeviceIDFactory.getInstance().getDeviceID(new
PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[j], tsDataType,
TSEncoding.PLAIN)),
System.nanoTime(),
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/memtable/PrimitiveMemTableTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/memtable/PrimitiveMemTableTest.java
index a6dbd1992d..0eca5e2a54 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/memtable/PrimitiveMemTableTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/memtable/PrimitiveMemTableTest.java
@@ -123,6 +123,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.write(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
dataSize - i - 1,
@@ -131,6 +132,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.write(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
i,
@@ -172,6 +174,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.write(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
i,
@@ -181,6 +184,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.write(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
schemaList,
i,
new Object[] {i, i});
@@ -192,6 +196,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.writeAlignedRow(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
schemaList,
i,
new Object[] {i, i});
@@ -199,6 +204,7 @@ public class PrimitiveMemTableTest {
Assert.assertEquals(5, memTable.getSeriesNumber());
memTable.writeAlignedRow(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[2], TSDataType.INT32,
TSEncoding.PLAIN)),
0,
@@ -220,6 +226,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.write(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
dataSize - i - 1,
@@ -228,6 +235,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.write(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
i,
@@ -274,6 +282,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.writeAlignedRow(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
dataSize - i - 1,
@@ -282,6 +291,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < dataSize; i++) {
memTable.writeAlignedRow(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(
new MeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
i,
@@ -327,6 +337,7 @@ public class PrimitiveMemTableTest {
for (TimeValuePair aRet : ret) {
memTable.write(
DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ Collections.emptyList(),
Collections.singletonList(new MeasurementSchema(sensorId, dataType,
encoding)),
aRet.getTimestamp(),
new Object[] {aRet.getValue().getValue()});
diff --git
a/server/src/test/java/org/apache/iotdb/db/writelog/recover/LogReplayerTest.java
b/server/src/test/java/org/apache/iotdb/db/writelog/recover/LogReplayerTest.java
index b63cd8fca6..618eb05af7 100644
---
a/server/src/test/java/org/apache/iotdb/db/writelog/recover/LogReplayerTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/writelog/recover/LogReplayerTest.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.writelog.recover;
+import org.apache.iotdb.db.conf.IoTDBConfig;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
import org.apache.iotdb.db.engine.memtable.IMemTable;
@@ -35,6 +36,7 @@ import
org.apache.iotdb.db.exception.metadata.MetadataException;
import org.apache.iotdb.db.exception.query.QueryProcessException;
import org.apache.iotdb.db.metadata.mnode.IMeasurementMNode;
import org.apache.iotdb.db.metadata.mnode.MeasurementMNode;
+import org.apache.iotdb.db.metadata.path.AlignedPath;
import org.apache.iotdb.db.metadata.path.MeasurementPath;
import org.apache.iotdb.db.metadata.path.PartialPath;
import org.apache.iotdb.db.qp.physical.crud.DeletePlan;
@@ -51,6 +53,7 @@ import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
import org.apache.iotdb.tsfile.read.TimeValuePair;
import org.apache.iotdb.tsfile.read.reader.IPointReader;
+import org.apache.iotdb.tsfile.utils.Binary;
import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
import org.junit.After;
@@ -62,6 +65,7 @@ import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.MappedByteBuffer;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Collections;
import java.util.List;
@@ -71,14 +75,26 @@ import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
public class LogReplayerTest {
+ IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+ boolean prevIsAutoCreateSchemaEnabled;
+ boolean prevIsEnablePartialInsert;
@Before
public void before() {
+ // set recover config, avoid creating deleted time series when recovering
wal
+ prevIsAutoCreateSchemaEnabled = config.isAutoCreateSchemaEnabled();
+ prevIsEnablePartialInsert = config.isEnablePartialInsert();
+ ;
+ config.setAutoCreateSchemaEnabled(false);
+ config.setEnablePartialInsert(true);
EnvironmentUtils.envSetUp();
}
@After
public void after() throws IOException, StorageEngineException {
+ // reset config
+ config.setAutoCreateSchemaEnabled(prevIsAutoCreateSchemaEnabled);
+ config.setEnablePartialInsert(prevIsEnablePartialInsert);
EnvironmentUtils.cleanEnv();
}
@@ -91,9 +107,11 @@ public class LogReplayerTest {
ModificationFile modFile = new ModificationFile(modF.getPath());
TsFileResource tsFileResource = new TsFileResource(tsFile);
IMemTable memTable = new PrimitiveMemTable();
+ CompressionType compressionType =
TSFileDescriptor.getInstance().getConfig().getCompressor();
IoTDB.metaManager.setStorageGroup(new PartialPath("root.sg"));
try {
+ // 1. set schema
for (int i = 0; i <= 5; i++) {
for (int j = 0; j <= 5; j++) {
IoTDB.metaManager.createTimeseries(
@@ -104,6 +122,25 @@ public class LogReplayerTest {
Collections.emptyMap());
}
}
+ IoTDB.metaManager.createAlignedTimeSeries(
+ new PartialPath("root.sg.device6"),
+ Arrays.asList("s1", "s2", "s3", "s4", "s5"),
+ Arrays.asList(
+ TSDataType.INT32,
+ TSDataType.INT64,
+ TSDataType.BOOLEAN,
+ TSDataType.FLOAT,
+ TSDataType.TEXT),
+ Arrays.asList(
+ TSEncoding.RLE, TSEncoding.RLE, TSEncoding.RLE, TSEncoding.RLE,
TSEncoding.PLAIN),
+ Arrays.asList(
+ compressionType, compressionType, compressionType,
compressionType, compressionType));
+
+ // 2. delete some timeseries
+ IoTDB.metaManager.deleteTimeseries(new
PartialPath("root.sg.device0.sensor2"));
+ IoTDB.metaManager.deleteTimeseries(new
PartialPath("root.sg.device0.sensor4"));
+ IoTDB.metaManager.deleteTimeseries(new
PartialPath("root.sg.device6.s1"));
+ IoTDB.metaManager.deleteTimeseries(new
PartialPath("root.sg.device6.s5"));
LogReplayer replayer =
new LogReplayer(
@@ -123,6 +160,13 @@ public class LogReplayerTest {
IoTDBDescriptor.getInstance().getConfig().getWalBufferSize() / 2);
return byteBuffers;
});
+ node.write(
+ new InsertRowPlan(
+ new PartialPath("root.sg.device0"),
+ 50,
+ "sensor4",
+ TSDataType.INT64,
+ String.valueOf(0)));
node.write(
new InsertRowPlan(
new PartialPath("root.sg.device0"),
@@ -146,7 +190,8 @@ public class LogReplayerTest {
TSDataType.INT64,
String.valueOf(i)));
}
- node.write(insertTablePlan());
+ node.write(insertTabletPlan());
+ node.write(insertAlignedTabletPlan());
DeletePlan deletePlan = new DeletePlan(0, 200, new
PartialPath("root.sg.device0.sensor0"));
node.write(deletePlan);
node.close();
@@ -185,6 +230,30 @@ public class LogReplayerTest {
}
assertFalse(iterator.hasNextTimeValuePair());
}
+ AlignedPath alignedfullPath =
+ new AlignedPath(
+ "root.sg.device6",
+ Arrays.asList("s1", "s2", "s3", "s4", "s5"),
+ Arrays.asList(
+ new MeasurementSchema("s1", TSDataType.INT32,
TSEncoding.RLE),
+ new MeasurementSchema("s2", TSDataType.INT64,
TSEncoding.RLE),
+ new MeasurementSchema("s3", TSDataType.BOOLEAN,
TSEncoding.RLE),
+ new MeasurementSchema("s4", TSDataType.FLOAT,
TSEncoding.RLE),
+ new MeasurementSchema("s5", TSDataType.TEXT,
TSEncoding.PLAIN)));
+ ReadOnlyMemChunk memChunk = memTable.query(alignedfullPath,
Long.MIN_VALUE, null);
+ IPointReader iterator = memChunk.getPointReader();
+ int time = 0;
+ while (iterator.hasNextTimeValuePair()) {
+ TimeValuePair timeValuePair = iterator.nextTimeValuePair();
+ assertEquals(time, timeValuePair.getTimestamp());
+ assertEquals(null, timeValuePair.getValue().getVector()[0]);
+ assertEquals(null, timeValuePair.getValue().getVector()[1]);
+ assertEquals(true,
timeValuePair.getValue().getVector()[2].getBoolean());
+ assertEquals(time, timeValuePair.getValue().getVector()[3].getFloat(),
0.00001);
+ assertEquals(null, timeValuePair.getValue().getVector()[4]);
+ time++;
+ }
+ assertEquals(100, time);
Modification[] mods = modFile.getModifications().toArray(new
Modification[0]);
assertEquals(1, mods.length);
@@ -194,6 +263,9 @@ public class LogReplayerTest {
assertEquals(2, tsFileResource.getStartTime("root.sg.device0"));
assertEquals(2, tsFileResource.getEndTime("root.sg.device0"));
+
+ assertEquals(0, tsFileResource.getStartTime("root.sg.device6"));
+ assertEquals(99, tsFileResource.getEndTime("root.sg.device6"));
for (int i = 1; i < 5; i++) {
assertEquals(i, tsFileResource.getStartTime("root.sg.device" + i));
assertEquals(i, tsFileResource.getEndTime("root.sg.device" + i));
@@ -211,14 +283,14 @@ public class LogReplayerTest {
TSEncoding.PLAIN,
CompressionType.UNCOMPRESSED,
Collections.emptyMap()));
- ReadOnlyMemChunk memChunk = memTable.query(fullPath, Long.MIN_VALUE,
null);
+ memChunk = memTable.query(fullPath, Long.MIN_VALUE, null);
// s0 has datatype boolean, but required INT64, will return null
if (i == 0) {
assertNull(memChunk);
} else {
- IPointReader iterator = memChunk.getPointReader();
+ iterator = memChunk.getPointReader();
iterator.hasNextTimeValuePair();
- for (int time = 0; time < 100; time++) {
+ for (time = 0; time < 100; time++) {
TimeValuePair timeValuePair = iterator.nextTimeValuePair();
assertEquals(time, timeValuePair.getTimestamp());
assertEquals(time, timeValuePair.getValue().getLong());
@@ -248,33 +320,43 @@ public class LogReplayerTest {
* @return
* @throws IllegalPathException
*/
- public InsertTabletPlan insertTablePlan() throws IllegalPathException {
- String[] measurements = new String[2];
- measurements[0] = "sensor0";
+ private InsertTabletPlan insertTabletPlan() throws IllegalPathException {
+ String[] measurements = new String[4];
+ measurements[0] = "sensor0"; // mismatch type
measurements[1] = "sensor1";
+ measurements[2] = "sensor2"; // have been deleted
+ measurements[3] = "sensor3";
List<Integer> dataTypes = new ArrayList<>();
dataTypes.add(TSDataType.BOOLEAN.ordinal());
dataTypes.add(TSDataType.INT64.ordinal());
+ dataTypes.add(TSDataType.INT64.ordinal());
+ dataTypes.add(TSDataType.INT64.ordinal());
String deviceId = "root.sg.device5";
- IMeasurementMNode[] mNodes = new IMeasurementMNode[2];
+ IMeasurementMNode[] mNodes = new IMeasurementMNode[4];
mNodes[0] = MeasurementMNode.getMeasurementMNode(null, "sensor0", null,
null);
mNodes[1] = MeasurementMNode.getMeasurementMNode(null, "sensor1", null,
null);
+ mNodes[2] = MeasurementMNode.getMeasurementMNode(null, "sensor2", null,
null);
+ mNodes[3] = MeasurementMNode.getMeasurementMNode(null, "sensor3", null,
null);
InsertTabletPlan insertTabletPlan =
new InsertTabletPlan(new PartialPath(deviceId), measurements,
dataTypes);
long[] times = new long[100];
- Object[] columns = new Object[2];
+ Object[] columns = new Object[4];
columns[0] = new boolean[100];
columns[1] = new long[100];
+ columns[2] = new long[100];
+ columns[3] = new long[100];
for (long r = 0; r < 100; r++) {
times[(int) r] = r;
((boolean[]) columns[0])[(int) r] = false;
((long[]) columns[1])[(int) r] = r;
+ ((long[]) columns[2])[(int) r] = r;
+ ((long[]) columns[3])[(int) r] = r;
}
insertTabletPlan.setTimes(times);
insertTabletPlan.setColumns(columns);
@@ -285,4 +367,60 @@ public class LogReplayerTest {
return insertTabletPlan;
}
+
+ /**
+ * insert tablet plan, time series expected datatype is INT64 s0 is set to
boolean, it will output
+ * null value s1 is set to INT64, it will output its value
+ *
+ * @return
+ * @throws IllegalPathException
+ */
+ private InsertTabletPlan insertAlignedTabletPlan() throws
IllegalPathException {
+ String deviceId = "root.sg.device6";
+
+ List<Integer> dataTypes =
+ Arrays.asList(
+ TSDataType.INT32.ordinal(), // deleted
+ TSDataType.BOOLEAN.ordinal(), // mismatch type
+ TSDataType.BOOLEAN.ordinal(),
+ TSDataType.FLOAT.ordinal(),
+ TSDataType.TEXT.ordinal()); // deleted
+
+ InsertTabletPlan insertTabletPlan =
+ new InsertTabletPlan(
+ new PartialPath(deviceId),
+ new String[] {"s1", "s2", "s3", "s4", "s5"},
+ dataTypes,
+ true);
+
+ long[] times = new long[100];
+ Object[] columns =
+ new Object[] {
+ new int[100], new boolean[100], new boolean[100], new float[100],
new Binary[100],
+ };
+ for (long r = 0; r < 100; r++) {
+ times[(int) r] = r;
+ ((int[]) columns[0])[(int) r] = (int) r;
+ ((boolean[]) columns[1])[(int) r] = true;
+ ((boolean[]) columns[2])[(int) r] = true;
+ ((float[]) columns[3])[(int) r] = r;
+ ((Binary[]) columns[4])[(int) r] = Binary.valueOf(r + "");
+ }
+
+ // BitMap[] bitMaps = new BitMap[dataTypes.size()];
+ // for (int i = 0; i < dataTypes.size(); i++) {
+ // if (bitMaps[i] == null) {
+ // bitMaps[i] = new BitMap(times.length);
+ // }
+ // // mark value of time=99 as null
+ // bitMaps[i].mark(99);
+ // }
+
+ insertTabletPlan.setTimes(times);
+ insertTabletPlan.setColumns(columns);
+ insertTabletPlan.setRowCount(times.length);
+ // insertTabletPlan.setBitMaps(bitMaps);
+
+ return insertTabletPlan;
+ }
}