This is an automated email from the ASF dual-hosted git repository.
jackietien pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 3d48a88 [IOTDB-1935] Aligned timeseries support deleting (#4454)
3d48a88 is described below
commit 3d48a88dc012165610cf7856e9308f1cf488fb97
Author: Haonan <[email protected]>
AuthorDate: Tue Nov 30 14:51:54 2021 +0800
[IOTDB-1935] Aligned timeseries support deleting (#4454)
---
.../iotdb/db/engine/memtable/AbstractMemTable.java | 29 +-
.../engine/memtable/AlignedWritableMemChunk.java | 27 +-
.../memtable/AlignedWritableMemChunkGroup.java | 27 ++
.../apache/iotdb/db/engine/memtable/IMemTable.java | 6 +-
.../db/engine/memtable/IWritableMemChunk.java | 3 -
.../db/engine/memtable/IWritableMemChunkGroup.java | 4 +
.../iotdb/db/engine/memtable/WritableMemChunk.java | 5 -
.../db/engine/memtable/WritableMemChunkGroup.java | 25 +
.../querycontext/AlignedReadOnlyMemChunk.java | 4 +-
.../db/engine/storagegroup/TsFileProcessor.java | 40 +-
.../apache/iotdb/db/metadata/path/AlignedPath.java | 32 +-
.../iotdb/db/metadata/path/MeasurementPath.java | 32 +-
.../apache/iotdb/db/metadata/path/PartialPath.java | 22 +-
.../db/utils/datastructure/AlignedTVList.java | 121 +++--
.../iotdb/db/utils/datastructure/TVList.java | 3 +-
.../db/engine/memtable/PrimitiveMemTableTest.java | 111 +++++
...gregationWithoutValueFilterWithDeletion2IT.java | 3 -
...ggregationWithoutValueFilterWithDeletionIT.java | 3 -
.../aligned/IoTDBDeleteTimeseriesIT.java | 215 +++++++++
.../db/integration/aligned/IoTDBDeletionIT.java | 519 +++++++++++++++++++++
.../aligned/IoTDBLastQueryWithDeletion2IT.java | 3 -
.../aligned/IoTDBLastQueryWithDeletionIT.java | 3 -
...DBLastQueryWithoutLastCacheWithDeletion2IT.java | 3 -
...TDBLastQueryWithoutLastCacheWithDeletionIT.java | 3 -
...oTDBRawQueryWithValueFilterWithDeletion2IT.java | 4 +-
...IoTDBRawQueryWithValueFilterWithDeletionIT.java | 6 +-
...BRawQueryWithoutValueFilterWithDeletion2IT.java | 4 +-
...DBRawQueryWithoutValueFilterWithDeletionIT.java | 12 +-
28 files changed, 1111 insertions(+), 158 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 e834529..ce79aff 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
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.engine.memtable;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.engine.modification.Modification;
import org.apache.iotdb.db.engine.querycontext.ReadOnlyMemChunk;
import org.apache.iotdb.db.exception.WriteProcessException;
import org.apache.iotdb.db.exception.query.QueryProcessException;
@@ -27,13 +28,12 @@ import org.apache.iotdb.db.qp.physical.crud.InsertRowPlan;
import org.apache.iotdb.db.qp.physical.crud.InsertTabletPlan;
import org.apache.iotdb.db.utils.MemUtils;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
-import org.apache.iotdb.tsfile.read.common.TimeRange;
+import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
import java.io.IOException;
import java.util.ArrayList;
import java.util.HashMap;
-import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
@@ -109,7 +109,7 @@ public abstract class AbstractMemTable implements IMemTable
{
deviceId,
k -> {
seriesNumber += schemaList.size();
- totalPointsNumThreshold += avgSeriesPointNumThreshold *
schemaList.size();
+ totalPointsNumThreshold += ((long) avgSeriesPointNumThreshold) *
schemaList.size();
return new AlignedWritableMemChunkGroup(schemaList);
});
for (IMeasurementSchema schema : schemaList) {
@@ -330,9 +330,9 @@ public abstract class AbstractMemTable implements IMemTable
{
@Override
public ReadOnlyMemChunk query(
- PartialPath fullPath, long ttlLowerBound, List<TimeRange> deletionList)
+ PartialPath fullPath, long ttlLowerBound, List<Pair<Modification,
IMemTable>> modsToMemtable)
throws IOException, QueryProcessException {
- return fullPath.getReadOnlyMemChunkFromMemTable(memTableMap, deletionList);
+ return fullPath.getReadOnlyMemChunkFromMemTable(this, modsToMemtable,
ttlLowerBound);
}
@SuppressWarnings("squid:S3776") // high Cognitive Complexity
@@ -343,24 +343,7 @@ public abstract class AbstractMemTable implements
IMemTable {
if (memChunkGroup == null) {
return;
}
-
- Iterator<Entry<String, IWritableMemChunk>> iter =
- memChunkGroup.getMemChunkMap().entrySet().iterator();
- while (iter.hasNext()) {
- Entry<String, IWritableMemChunk> entry = iter.next();
- IWritableMemChunk chunk = entry.getValue();
- // the key is measurement rather than component of multiMeasurement
- PartialPath fullPath = devicePath.concatNode(entry.getKey());
- if (originalPath.matchFullPath(fullPath)) {
- // matchFullPath ensures this branch could work on delete data of
unary or multi measurement
- // and delete timeseries or aligned timeseries
- if (startTimestamp == Long.MIN_VALUE && endTimestamp ==
Long.MAX_VALUE) {
- iter.remove();
- }
- int deletedPointsNumber = chunk.delete(startTimestamp, endTimestamp);
- totalPointsNum -= deletedPointsNumber;
- }
- }
+ totalPointsNum -= memChunkGroup.delete(originalPath, devicePath,
startTimestamp, endTimestamp);
}
@Override
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 001a3c2..6cfed35 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
@@ -25,6 +25,7 @@ import
org.apache.iotdb.tsfile.exception.write.UnSupportedDataTypeException;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.utils.Binary;
import org.apache.iotdb.tsfile.utils.BitMap;
+import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.chunk.AlignedChunkWriterImpl;
import org.apache.iotdb.tsfile.write.chunk.IChunkWriter;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
@@ -36,6 +37,7 @@ import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.Set;
public class AlignedWritableMemChunk implements IWritableMemChunk {
@@ -56,6 +58,10 @@ public class AlignedWritableMemChunk implements
IWritableMemChunk {
this.list = TVListAllocator.getInstance().allocate(dataTypeList);
}
+ public Set<String> getAllMeasurements() {
+ return measurementIndexMap.keySet();
+ }
+
public boolean containsMeasurement(String measurementId) {
return measurementIndexMap.containsKey(measurementId);
}
@@ -241,10 +247,23 @@ public class AlignedWritableMemChunk implements
IWritableMemChunk {
return list.delete(lowerBound, upperBound);
}
- @Override
- // TODO: THIS METHOLD IS FOR DELETING ONE COLUMN OF A VECTOR
- public int delete(long lowerBound, long upperBound, String measurementId) {
- return 0;
+ public Pair<Integer, Boolean> deleteDataFromAColumn(
+ long lowerBound, long upperBound, String measurementId) {
+ return list.delete(lowerBound, upperBound,
measurementIndexMap.get(measurementId));
+ }
+
+ public void removeColumns(List<String> measurements) {
+ List<IMeasurementSchema> schemasToBeRemoved = new ArrayList<>();
+ for (String measurement : measurements) {
+
schemasToBeRemoved.add(schemaList.get(measurementIndexMap.get(measurement)));
+ }
+ for (IMeasurementSchema schema : schemasToBeRemoved) {
+ schemaList.remove(schema);
+ }
+ measurementIndexMap.clear();
+ for (int i = 0; i < schemaList.size(); i++) {
+ measurementIndexMap.put(schemaList.get(i).getMeasurementId(), i);
+ }
}
@Override
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 0a62bd6..b19a59d 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
@@ -19,12 +19,16 @@
package org.apache.iotdb.db.engine.memtable;
+import org.apache.iotdb.db.metadata.path.PartialPath;
import org.apache.iotdb.tsfile.utils.BitMap;
+import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
+import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
+import java.util.Set;
public class AlignedWritableMemChunkGroup implements IWritableMemChunkGroup {
@@ -71,6 +75,29 @@ public class AlignedWritableMemChunkGroup implements
IWritableMemChunkGroup {
}
@Override
+ public int delete(
+ PartialPath originalPath, PartialPath devicePath, long startTimestamp,
long endTimestamp) {
+ int deletedPointsNumber = 0;
+ Set<String> measurements = memChunk.getAllMeasurements();
+ List<String> columnsToBeRemoved = new ArrayList<>();
+ for (String measurement : measurements) {
+ PartialPath fullPath = devicePath.concatNode(measurement);
+ if (originalPath.matchFullPath(fullPath)) {
+ Pair<Integer, Boolean> deleteInfo =
+ memChunk.deleteDataFromAColumn(startTimestamp, endTimestamp,
measurement);
+ deletedPointsNumber += deleteInfo.left;
+ if (Boolean.TRUE.equals(deleteInfo.right)) {
+ columnsToBeRemoved.add(measurement);
+ }
+ }
+ }
+ if (!columnsToBeRemoved.isEmpty()) {
+ memChunk.removeColumns(columnsToBeRemoved);
+ }
+ return deletedPointsNumber;
+ }
+
+ @Override
public long getCurrentChunkPointNum(String measurement) {
return memChunk.count();
}
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 e05b1be..fa228f4 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
@@ -18,6 +18,7 @@
*/
package org.apache.iotdb.db.engine.memtable;
+import org.apache.iotdb.db.engine.modification.Modification;
import org.apache.iotdb.db.engine.querycontext.ReadOnlyMemChunk;
import org.apache.iotdb.db.exception.WriteProcessException;
import org.apache.iotdb.db.exception.metadata.MetadataException;
@@ -25,7 +26,7 @@ import
org.apache.iotdb.db.exception.query.QueryProcessException;
import org.apache.iotdb.db.metadata.path.PartialPath;
import org.apache.iotdb.db.qp.physical.crud.InsertRowPlan;
import org.apache.iotdb.db.qp.physical.crud.InsertTabletPlan;
-import org.apache.iotdb.tsfile.read.common.TimeRange;
+import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
import java.io.IOException;
@@ -107,7 +108,8 @@ public interface IMemTable {
void insertAlignedTablet(InsertTabletPlan insertTabletPlan, int start, int
end)
throws WriteProcessException;
- ReadOnlyMemChunk query(PartialPath fullPath, long ttlLowerBound,
List<TimeRange> deletionList)
+ ReadOnlyMemChunk query(
+ PartialPath fullPath, long ttlLowerBound, List<Pair<Modification,
IMemTable>> modsToMemtable)
throws IOException, QueryProcessException, MetadataException;
/** putBack all the memory resources. */
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 c2c1781..3061c30 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
@@ -124,9 +124,6 @@ public interface IWritableMemChunk {
/** @return how many points are deleted */
int delete(long lowerBound, long upperBound);
- // For delete one column in the vector
- int delete(long lowerBound, long upperBound, String measurementId);
-
IChunkWriter createIChunkWriter();
void encode(IChunkWriter chunkWriter);
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 85586aa..00bcf7c 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
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.engine.memtable;
+import org.apache.iotdb.db.metadata.path.PartialPath;
import org.apache.iotdb.tsfile.utils.BitMap;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
@@ -45,5 +46,8 @@ public interface IWritableMemChunkGroup {
Map<String, IWritableMemChunk> getMemChunkMap();
+ int delete(
+ PartialPath originalPath, PartialPath devicePath, long startTimestamp,
long endTimestamp);
+
long getCurrentChunkPointNum(String measurement);
}
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 e5f288c..df107ce 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
@@ -265,11 +265,6 @@ public class WritableMemChunk implements IWritableMemChunk
{
}
@Override
- public int delete(long lowerBound, long upperBound, String measurementId) {
- throw new UnSupportedDataTypeException(UNSUPPORTED_TYPE +
schema.getType());
- }
-
- @Override
public IChunkWriter createIChunkWriter() {
return new ChunkWriterImpl(schema);
}
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 fd6c408..578a346 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
@@ -19,12 +19,15 @@
package org.apache.iotdb.db.engine.memtable;
+import org.apache.iotdb.db.metadata.path.PartialPath;
import org.apache.iotdb.tsfile.utils.BitMap;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
import java.util.HashMap;
+import java.util.Iterator;
import java.util.List;
import java.util.Map;
+import java.util.Map.Entry;
public class WritableMemChunkGroup implements IWritableMemChunkGroup {
@@ -106,6 +109,28 @@ public class WritableMemChunkGroup implements
IWritableMemChunkGroup {
}
@Override
+ public int delete(
+ PartialPath originalPath, PartialPath devicePath, long startTimestamp,
long endTimestamp) {
+ int deletedPointsNumber = 0;
+ Iterator<Entry<String, IWritableMemChunk>> iter =
memChunkMap.entrySet().iterator();
+ while (iter.hasNext()) {
+ Entry<String, IWritableMemChunk> entry = iter.next();
+ IWritableMemChunk chunk = entry.getValue();
+ // the key is measurement rather than component of multiMeasurement
+ PartialPath fullPath = devicePath.concatNode(entry.getKey());
+ if (originalPath.matchFullPath(fullPath)) {
+ // matchFullPath ensures this branch could work on delete data of
unary or multi measurement
+ // and delete timeseries
+ if (startTimestamp == Long.MIN_VALUE && endTimestamp ==
Long.MAX_VALUE) {
+ iter.remove();
+ }
+ deletedPointsNumber += chunk.delete(startTimestamp, endTimestamp);
+ }
+ }
+ return deletedPointsNumber;
+ }
+
+ @Override
public long getCurrentChunkPointNum(String measurement) {
return memChunkMap.get(measurement).count();
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/querycontext/AlignedReadOnlyMemChunk.java
b/server/src/main/java/org/apache/iotdb/db/engine/querycontext/AlignedReadOnlyMemChunk.java
index 2f18ef7..1e058a1 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/querycontext/AlignedReadOnlyMemChunk.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/querycontext/AlignedReadOnlyMemChunk.java
@@ -45,7 +45,7 @@ import java.util.List;
public class AlignedReadOnlyMemChunk extends ReadOnlyMemChunk {
// deletion list for this chunk
- private final List<TimeRange> deletionList;
+ private final List<List<TimeRange>> deletionList;
private String measurementUid;
private TSDataType dataType;
@@ -68,7 +68,7 @@ public class AlignedReadOnlyMemChunk extends ReadOnlyMemChunk
{
* @param deletionList The timeRange of deletionList
*/
public AlignedReadOnlyMemChunk(
- IMeasurementSchema schema, TVList tvList, int size, List<TimeRange>
deletionList)
+ IMeasurementSchema schema, TVList tvList, int size,
List<List<TimeRange>> deletionList)
throws IOException, QueryProcessException {
super();
this.measurementUid = schema.getMeasurementId();
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
index 5908597..e52a43d 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
@@ -60,7 +60,6 @@ import org.apache.iotdb.service.rpc.thrift.TSStatus;
import org.apache.iotdb.tsfile.file.metadata.ChunkMetadata;
import org.apache.iotdb.tsfile.file.metadata.IChunkMetadata;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
-import org.apache.iotdb.tsfile.read.common.TimeRange;
import org.apache.iotdb.tsfile.utils.Binary;
import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter;
@@ -1227,41 +1226,6 @@ public class TsFileProcessor {
return storageGroupName;
}
- /** get modifications from a memtable */
- private List<Modification> getModificationsForMemtable(IMemTable memTable) {
- List<Modification> modifications = new ArrayList<>();
- boolean foundMemtable = false;
- for (Pair<Modification, IMemTable> entry : modsToMemtable) {
- if (foundMemtable || entry.right.equals(memTable)) {
- modifications.add(entry.left);
- foundMemtable = true;
- }
- }
- return modifications;
- }
-
- /**
- * construct a deletion list from a memtable
- *
- * @param memTable memtable
- * @param timeLowerBound time water mark
- */
- private List<TimeRange> constructDeletionList(
- IMemTable memTable, PartialPath fullPath, long timeLowerBound) {
- List<TimeRange> deletionList = new ArrayList<>();
- deletionList.add(new TimeRange(Long.MIN_VALUE, timeLowerBound));
- for (Modification modification : getModificationsForMemtable(memTable)) {
- if (modification instanceof Deletion) {
- Deletion deletion = (Deletion) modification;
- if (deletion.getPath().matchFullPath(fullPath) &&
deletion.getEndTime() > timeLowerBound) {
- long lowerBound = Math.max(deletion.getStartTime(), timeLowerBound);
- deletionList.add(new TimeRange(lowerBound, deletion.getEndTime()));
- }
- }
- }
- return TimeRange.sortAndMerge(deletionList);
- }
-
/**
* get the chunk(s) in the memtable (one from work memtable and the other
ones in flushing
* memtables and then compact them into one TimeValuePairSorter). Then get
the related
@@ -1284,10 +1248,8 @@ public class TsFileProcessor {
if (flushingMemTable.isSignalMemTable()) {
continue;
}
- List<TimeRange> deletionList =
- constructDeletionList(flushingMemTable, fullPath,
context.getQueryTimeLowerBound());
ReadOnlyMemChunk memChunk =
- flushingMemTable.query(fullPath, context.getQueryTimeLowerBound(),
deletionList);
+ flushingMemTable.query(fullPath, context.getQueryTimeLowerBound(),
modsToMemtable);
if (memChunk != null) {
readOnlyMemChunks.add(memChunk);
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/path/AlignedPath.java
b/server/src/main/java/org/apache/iotdb/db/metadata/path/AlignedPath.java
index 7a8ee70..8d9124b 100644
--- a/server/src/main/java/org/apache/iotdb/db/metadata/path/AlignedPath.java
+++ b/server/src/main/java/org/apache/iotdb/db/metadata/path/AlignedPath.java
@@ -21,7 +21,9 @@ package org.apache.iotdb.db.metadata.path;
import org.apache.iotdb.db.engine.memtable.AlignedWritableMemChunk;
import org.apache.iotdb.db.engine.memtable.AlignedWritableMemChunkGroup;
+import org.apache.iotdb.db.engine.memtable.IMemTable;
import org.apache.iotdb.db.engine.memtable.IWritableMemChunkGroup;
+import org.apache.iotdb.db.engine.modification.Deletion;
import org.apache.iotdb.db.engine.modification.Modification;
import org.apache.iotdb.db.engine.modification.ModificationFile;
import org.apache.iotdb.db.engine.querycontext.AlignedReadOnlyMemChunk;
@@ -47,6 +49,7 @@ import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics;
import org.apache.iotdb.tsfile.read.common.TimeRange;
import org.apache.iotdb.tsfile.read.filter.basic.Filter;
+import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
import org.apache.iotdb.tsfile.write.schema.VectorMeasurementSchema;
import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter;
@@ -357,8 +360,9 @@ public class AlignedPath extends PartialPath {
@Override
public ReadOnlyMemChunk getReadOnlyMemChunkFromMemTable(
- Map<String, IWritableMemChunkGroup> memTableMap, List<TimeRange>
deletionList)
+ IMemTable memTable, List<Pair<Modification, IMemTable>> modsToMemtable,
long timeLowerBound)
throws QueryProcessException, IOException {
+ Map<String, IWritableMemChunkGroup> memTableMap =
memTable.getMemTableMap();
// check If memtable contains this path
if (!memTableMap.containsKey(getDevice())) {
return null;
@@ -378,10 +382,36 @@ public class AlignedPath extends PartialPath {
// get sorted tv list is synchronized so different query can get right
sorted list reference
TVList alignedTvListCopy =
alignedMemChunk.getSortedTvListForQuery(schemaList);
int curSize = alignedTvListCopy.size();
+ List<List<TimeRange>> deletionList = null;
+ if (modsToMemtable != null) {
+ deletionList = constructDeletionList(memTable, modsToMemtable,
timeLowerBound);
+ }
return new AlignedReadOnlyMemChunk(
getMeasurementSchema(), alignedTvListCopy, curSize, deletionList);
}
+ private List<List<TimeRange>> constructDeletionList(
+ IMemTable memTable, List<Pair<Modification, IMemTable>> modsToMemtable,
long timeLowerBound) {
+ List<List<TimeRange>> deletionList = new ArrayList<>();
+ for (String measurement : measurementList) {
+ List<TimeRange> columnDeletionList = new ArrayList<>();
+ columnDeletionList.add(new TimeRange(Long.MIN_VALUE, timeLowerBound));
+ for (Modification modification : getModificationsForMemtable(memTable,
modsToMemtable)) {
+ if (modification instanceof Deletion) {
+ Deletion deletion = (Deletion) modification;
+ PartialPath fullPath = this.concatNode(measurement);
+ if (deletion.getPath().matchFullPath(fullPath)
+ && deletion.getEndTime() > timeLowerBound) {
+ long lowerBound = Math.max(deletion.getStartTime(),
timeLowerBound);
+ columnDeletionList.add(new TimeRange(lowerBound,
deletion.getEndTime()));
+ }
+ }
+ }
+ deletionList.add(TimeRange.sortAndMerge(columnDeletionList));
+ }
+ return deletionList;
+ }
+
@Override
public List<IChunkMetadata> getVisibleMetadataListFromWriter(
RestorableTsFileIOWriter writer, TsFileResource tsFileResource,
QueryContext context) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/path/MeasurementPath.java
b/server/src/main/java/org/apache/iotdb/db/metadata/path/MeasurementPath.java
index 9d4277a..c491907 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/path/MeasurementPath.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/path/MeasurementPath.java
@@ -19,8 +19,10 @@
package org.apache.iotdb.db.metadata.path;
import org.apache.iotdb.db.conf.IoTDBConstant;
+import org.apache.iotdb.db.engine.memtable.IMemTable;
import org.apache.iotdb.db.engine.memtable.IWritableMemChunk;
import org.apache.iotdb.db.engine.memtable.IWritableMemChunkGroup;
+import org.apache.iotdb.db.engine.modification.Deletion;
import org.apache.iotdb.db.engine.modification.Modification;
import org.apache.iotdb.db.engine.modification.ModificationFile;
import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
@@ -41,6 +43,7 @@ import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics;
import org.apache.iotdb.tsfile.read.common.TimeRange;
import org.apache.iotdb.tsfile.read.filter.basic.Filter;
+import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
import org.apache.iotdb.tsfile.write.schema.UnaryMeasurementSchema;
import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter;
@@ -247,8 +250,9 @@ public class MeasurementPath extends PartialPath {
@Override
public ReadOnlyMemChunk getReadOnlyMemChunkFromMemTable(
- Map<String, IWritableMemChunkGroup> memTableMap, List<TimeRange>
deletionList)
+ IMemTable memTable, List<Pair<Modification, IMemTable>> modsToMemtable,
long timeLowerBound)
throws QueryProcessException, IOException {
+ Map<String, IWritableMemChunkGroup> memTableMap =
memTable.getMemTableMap();
// check If Memtable Contains this path
if (!memTableMap.containsKey(getDevice())
|| !memTableMap.get(getDevice()).contains(getMeasurement())) {
@@ -259,6 +263,10 @@ public class MeasurementPath extends PartialPath {
// get sorted tv list is synchronized so different query can get right
sorted list reference
TVList chunkCopy = memChunk.getSortedTvListForQuery();
int curSize = chunkCopy.size();
+ List<TimeRange> deletionList = null;
+ if (modsToMemtable != null) {
+ deletionList = constructDeletionList(memTable, modsToMemtable,
timeLowerBound);
+ }
return new ReadOnlyMemChunk(
getMeasurement(),
measurementSchema.getType(),
@@ -269,6 +277,28 @@ public class MeasurementPath extends PartialPath {
deletionList);
}
+ /**
+ * construct a deletion list from a memtable.
+ *
+ * @param memTable memtable
+ * @param timeLowerBound time water mark
+ */
+ private List<TimeRange> constructDeletionList(
+ IMemTable memTable, List<Pair<Modification, IMemTable>> modsToMemtable,
long timeLowerBound) {
+ List<TimeRange> deletionList = new ArrayList<>();
+ deletionList.add(new TimeRange(Long.MIN_VALUE, timeLowerBound));
+ for (Modification modification : getModificationsForMemtable(memTable,
modsToMemtable)) {
+ if (modification instanceof Deletion) {
+ Deletion deletion = (Deletion) modification;
+ if (deletion.getPath().matchFullPath(this) && deletion.getEndTime() >
timeLowerBound) {
+ long lowerBound = Math.max(deletion.getStartTime(), timeLowerBound);
+ deletionList.add(new TimeRange(lowerBound, deletion.getEndTime()));
+ }
+ }
+ }
+ return TimeRange.sortAndMerge(deletionList);
+ }
+
@Override
public MeasurementPath clone() {
MeasurementPath newMeasurementPath = null;
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/path/PartialPath.java
b/server/src/main/java/org/apache/iotdb/db/metadata/path/PartialPath.java
index ae42486..405f214 100644
--- a/server/src/main/java/org/apache/iotdb/db/metadata/path/PartialPath.java
+++ b/server/src/main/java/org/apache/iotdb/db/metadata/path/PartialPath.java
@@ -18,7 +18,8 @@
*/
package org.apache.iotdb.db.metadata.path;
-import org.apache.iotdb.db.engine.memtable.IWritableMemChunkGroup;
+import org.apache.iotdb.db.engine.memtable.IMemTable;
+import org.apache.iotdb.db.engine.modification.Modification;
import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
import org.apache.iotdb.db.engine.querycontext.ReadOnlyMemChunk;
import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
@@ -35,8 +36,8 @@ import org.apache.iotdb.tsfile.common.constant.TsFileConstant;
import org.apache.iotdb.tsfile.file.metadata.IChunkMetadata;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.read.common.Path;
-import org.apache.iotdb.tsfile.read.common.TimeRange;
import org.apache.iotdb.tsfile.read.filter.basic.Filter;
+import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter;
@@ -48,7 +49,6 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
-import java.util.Map;
import java.util.Set;
import java.util.regex.Pattern;
@@ -423,11 +423,25 @@ public class PartialPath extends Path implements
Comparable<Path>, Cloneable {
* @return ReadOnlyMemChunk
*/
public ReadOnlyMemChunk getReadOnlyMemChunkFromMemTable(
- Map<String, IWritableMemChunkGroup> memTableMap, List<TimeRange>
deletionList)
+ IMemTable memTable, List<Pair<Modification, IMemTable>> modsToMemtable,
long timeLowerBound)
throws QueryProcessException, IOException {
throw new UnsupportedOperationException("Should call exact sub class!");
}
+ /** get modifications from a memtable. */
+ protected List<Modification> getModificationsForMemtable(
+ IMemTable memTable, List<Pair<Modification, IMemTable>> modsToMemtable) {
+ List<Modification> modifications = new ArrayList<>();
+ boolean foundMemtable = false;
+ for (Pair<Modification, IMemTable> entry : modsToMemtable) {
+ if (foundMemtable || entry.right.equals(memTable)) {
+ modifications.add(entry.left);
+ foundMemtable = true;
+ }
+ }
+ return modifications;
+ }
+
@Override
public PartialPath clone() {
return new PartialPath(this.getNodes().clone());
diff --git
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java
index 1c9f6ed..120c9f8 100644
---
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java
+++
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java
@@ -29,6 +29,7 @@ import org.apache.iotdb.tsfile.read.common.TimeRange;
import org.apache.iotdb.tsfile.read.reader.IPointReader;
import org.apache.iotdb.tsfile.utils.Binary;
import org.apache.iotdb.tsfile.utils.BitMap;
+import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
import java.util.ArrayList;
@@ -176,10 +177,11 @@ public class AlignedTVList extends TVList {
TsPrimitiveType[] vector = new TsPrimitiveType[values.size()];
for (int columnIndex = 0; columnIndex < values.size(); columnIndex++) {
List<Object> columnValues = values.get(columnIndex);
- if (columnValues == null
- || bitMaps != null
- && bitMaps.get(columnIndex) != null
- && isValueMarked(valueIndex, columnIndex)) {
+ if (validIndexesForTimeDuplicatedRows == null
+ && (columnValues == null
+ || bitMaps != null
+ && bitMaps.get(columnIndex) != null
+ && isValueMarked(valueIndex, columnIndex))) {
continue;
}
if (validIndexesForTimeDuplicatedRows != null) {
@@ -267,8 +269,8 @@ public class AlignedTVList extends TVList {
public void extendColumn(TSDataType dataType) {
if (bitMaps == null) {
- bitMaps = new ArrayList<>(values.size() + 1);
- for (int i = 0; i < values.size() + 1; i++) {
+ bitMaps = new ArrayList<>(values.size());
+ for (int i = 0; i < values.size(); i++) {
bitMaps.add(null);
}
}
@@ -308,8 +310,7 @@ public class AlignedTVList extends TVList {
}
columnBitMaps.add(bitMap);
}
- // values.size() is the index of column
- this.bitMaps.set(values.size(), columnBitMaps);
+ this.bitMaps.add(columnBitMaps);
this.values.add(columnValue);
this.dataTypes.add(dataType);
}
@@ -430,35 +431,48 @@ public class AlignedTVList extends TVList {
@Override
public int delete(long lowerBound, long upperBound) {
- int newSize = 0;
- minTime = Long.MAX_VALUE;
+ int deletedNumber = 0;
+ for (int i = 0; i < dataTypes.size(); i++) {
+ deletedNumber += delete(lowerBound, upperBound, i).left;
+ }
+ return deletedNumber;
+ }
+
+ /**
+ * Delete points in a specific column.
+ *
+ * @param lowerBound deletion lower bound
+ * @param upperBound deletion upper bound
+ * @param columnIndex column index to be deleted
+ * @return Delete info pair. Left: deletedNumber int; right: ifDeleteColumn
boolean
+ */
+ public Pair<Integer, Boolean> delete(long lowerBound, long upperBound, int
columnIndex) {
+ int deletedNumber = 0;
+ boolean deleteColumn = true;
for (int i = 0; i < size; i++) {
long time = getTime(i);
- if (time < lowerBound || time > upperBound) {
- set(i, newSize++);
- minTime = Math.min(time, minTime);
+ if (time >= lowerBound && time <= upperBound) {
+ int originRowIndex = getValueIndex(i);
+ int arrayIndex = originRowIndex / ARRAY_SIZE;
+ int elementIndex = originRowIndex % ARRAY_SIZE;
+ markNullValue(columnIndex, arrayIndex, elementIndex);
+ deletedNumber++;
+ } else {
+ deleteColumn = false;
}
}
- int deletedNumber = size - newSize;
- size = newSize;
- // release primitive arrays that are empty
- int newArrayNum = newSize / ARRAY_SIZE;
- if (newSize % ARRAY_SIZE != 0) {
- newArrayNum++;
- }
- for (int releaseIdx = newArrayNum; releaseIdx < timestamps.size();
releaseIdx++) {
- releaseLastTimeArray();
- releaseLastValueArray();
+ if (deleteColumn) {
+ dataTypes.remove(columnIndex);
+ for (Object array : values.get(columnIndex)) {
+ PrimitiveArrayManager.release(array);
+ }
+ values.remove(columnIndex);
+ bitMaps.remove(columnIndex);
}
- return deletedNumber * getTsDataTypes().size();
+ return new Pair<>(deletedNumber, deleteColumn);
}
- // TODO: THIS METHOLD IS FOR DELETING ONE COLUMN OF A VECTOR
- public int delete(long lowerBound, long upperBound, int columnIndex) {
- throw new UnsupportedOperationException(ERR_DATATYPE_NOT_CONSISTENT);
- }
-
- protected void set(int index, long timestamp, int value) {
+ private void set(int index, long timestamp, int value) {
int arrayIndex = index / ARRAY_SIZE;
int elementIndex = index % ARRAY_SIZE;
timestamps.get(arrayIndex)[elementIndex] = timestamp;
@@ -816,7 +830,7 @@ public class AlignedTVList extends TVList {
if (bitMaps.get(columnIndex) == null) {
List<BitMap> columnBitMaps = new ArrayList<>();
for (int i = 0; i < values.get(columnIndex).size(); i++) {
- columnBitMaps.add(null);
+ columnBitMaps.add(new BitMap(ARRAY_SIZE));
}
bitMaps.set(columnIndex, columnBitMaps);
}
@@ -879,22 +893,35 @@ public class AlignedTVList extends TVList {
}
public IPointReader getAlignedIterator(
- int floatPrecision, List<TSEncoding> encodingList, int size,
List<TimeRange> deletionList) {
+ int floatPrecision,
+ List<TSEncoding> encodingList,
+ int size,
+ List<List<TimeRange>> deletionList) {
return new AlignedIte(floatPrecision, encodingList, size, deletionList);
}
private class AlignedIte extends Ite {
private List<TSEncoding> encodingList;
+ private int[] deleteCursors;
+ /** this field is effective only in the AlignedTvlist in a
AlignedRealOnlyMemChunk. */
+ private List<List<TimeRange>> deletionList;
public AlignedIte() {
super();
}
public AlignedIte(
- int floatPrecision, List<TSEncoding> encodingList, int size,
List<TimeRange> deletionList) {
- super(floatPrecision, null, size, deletionList);
+ int floatPrecision,
+ List<TSEncoding> encodingList,
+ int size,
+ List<List<TimeRange>> deletionList) {
+ super(floatPrecision, null, size, null);
this.encodingList = encodingList;
+ this.deletionList = deletionList;
+ if (deletionList != null) {
+ deleteCursors = new int[deletionList.size()];
+ }
}
@Override
@@ -906,7 +933,7 @@ public class AlignedTVList extends TVList {
List<Integer> timeDuplicatedAlignedRowIndexList = null;
while (cur < iteSize) {
long time = getTime(cur);
- if (isPointDeleted(time) || (cur + 1 < size() && (time == getTime(cur
+ 1)))) {
+ if (cur + 1 < size() && (time == getTime(cur + 1))) {
if (timeDuplicatedAlignedRowIndexList == null) {
timeDuplicatedAlignedRowIndexList = new ArrayList<>();
timeDuplicatedAlignedRowIndexList.add(getValueIndex(cur));
@@ -925,6 +952,9 @@ public class AlignedTVList extends TVList {
tvPair = getTimeValuePair(cur, time, floatPrecision, encodingList);
}
cur++;
+ if (deletePointsInDeletionList(time, tvPair)) {
+ continue;
+ }
if (tvPair.getValue() != null) {
cachedTimeValuePair = tvPair;
hasCachedPair = true;
@@ -934,5 +964,26 @@ public class AlignedTVList extends TVList {
return false;
}
+
+ private boolean deletePointsInDeletionList(long timestamp, TimeValuePair
tvPair) {
+ if (deletionList == null) {
+ return false;
+ }
+ boolean deletedAll = true;
+ for (int i = 0; i < deleteCursors.length; i++) {
+ while (deletionList.get(i) != null && deleteCursors[i] <
deletionList.get(i).size()) {
+ if (deletionList.get(i).get(deleteCursors[i]).contains(timestamp)) {
+ tvPair.getValue().getVector()[i] = null;
+ break;
+ } else if (deletionList.get(i).get(deleteCursors[i]).getMax() <
timestamp) {
+ deleteCursors[i]++;
+ } else {
+ deletedAll = false;
+ break;
+ }
+ }
+ }
+ return deletedAll;
+ }
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java
index 2d895c4..0b2cf6f 100644
--- a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java
+++ b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java
@@ -266,7 +266,8 @@ public abstract class TVList {
if (newSize % ARRAY_SIZE != 0) {
newArrayNum++;
}
- for (int releaseIdx = newArrayNum; releaseIdx < timestamps.size();
releaseIdx++) {
+ int oldArrayNum = timestamps.size();
+ for (int releaseIdx = newArrayNum; releaseIdx < oldArrayNum; releaseIdx++)
{
releaseLastTimeArray();
releaseLastValueArray();
}
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 a6f61e8..4232971 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
@@ -18,6 +18,8 @@
*/
package org.apache.iotdb.db.engine.memtable;
+import org.apache.iotdb.db.engine.modification.Deletion;
+import org.apache.iotdb.db.engine.modification.Modification;
import org.apache.iotdb.db.engine.querycontext.ReadOnlyMemChunk;
import org.apache.iotdb.db.exception.metadata.IllegalPathException;
import org.apache.iotdb.db.exception.metadata.MetadataException;
@@ -37,6 +39,7 @@ 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.utils.Pair;
import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
import org.apache.iotdb.tsfile.write.schema.UnaryMeasurementSchema;
@@ -196,6 +199,114 @@ public class PrimitiveMemTableTest {
Assert.assertEquals(6, memTable.getSeriesNumber());
}
+ @Test
+ public void queryWithDeletionTest() throws IOException,
QueryProcessException, MetadataException {
+ IMemTable memTable = new PrimitiveMemTable();
+ int count = 10;
+ String deviceId = "d1";
+ String[] measurementId = new String[count];
+ for (int i = 0; i < measurementId.length; i++) {
+ measurementId[i] = "s" + i;
+ }
+
+ int dataSize = 10000;
+ for (int i = 0; i < dataSize; i++) {
+ memTable.write(
+ deviceId,
+ Collections.singletonList(
+ new UnaryMeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
+ dataSize - i - 1,
+ new Object[] {i + 10});
+ }
+ for (int i = 0; i < dataSize; i++) {
+ memTable.write(
+ deviceId,
+ Collections.singletonList(
+ new UnaryMeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
+ i,
+ new Object[] {i});
+ }
+ MeasurementPath fullPath =
+ new MeasurementPath(
+ deviceId,
+ measurementId[0],
+ new UnaryMeasurementSchema(
+ measurementId[0],
+ TSDataType.INT32,
+ TSEncoding.RLE,
+ CompressionType.UNCOMPRESSED,
+ Collections.emptyMap()));
+ List<Pair<Modification, IMemTable>> modsToMemtable = new ArrayList<>();
+ Modification deletion =
+ new Deletion(new PartialPath(deviceId, measurementId[0]),
Long.MAX_VALUE, 10, dataSize);
+ modsToMemtable.add(new Pair<>(deletion, memTable));
+ ReadOnlyMemChunk memChunk = memTable.query(fullPath, Long.MIN_VALUE,
modsToMemtable);
+ IPointReader iterator = memChunk.getPointReader();
+ int cnt = 0;
+ while (iterator.hasNextTimeValuePair()) {
+ TimeValuePair timeValuePair = iterator.nextTimeValuePair();
+ Assert.assertEquals(cnt, timeValuePair.getTimestamp());
+ Assert.assertEquals(cnt, timeValuePair.getValue().getValue());
+ cnt++;
+ }
+ Assert.assertEquals(10, cnt);
+ }
+
+ @Test
+ public void queryAlignChuckWithDeletionTest()
+ throws IOException, QueryProcessException, MetadataException {
+ IMemTable memTable = new PrimitiveMemTable();
+ int count = 10;
+ String deviceId = "d1";
+ String[] measurementId = new String[count];
+ for (int i = 0; i < measurementId.length; i++) {
+ measurementId[i] = "s" + i;
+ }
+
+ int dataSize = 10000;
+ for (int i = 0; i < dataSize; i++) {
+ memTable.writeAlignedRow(
+ deviceId,
+ Collections.singletonList(
+ new UnaryMeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
+ dataSize - i - 1,
+ new Object[] {i + 10});
+ }
+ for (int i = 0; i < dataSize; i++) {
+ memTable.writeAlignedRow(
+ deviceId,
+ Collections.singletonList(
+ new UnaryMeasurementSchema(measurementId[0], TSDataType.INT32,
TSEncoding.PLAIN)),
+ i,
+ new Object[] {i});
+ }
+ AlignedPath fullPath =
+ new AlignedPath(
+ deviceId,
+ Collections.singletonList(measurementId[0]),
+ Collections.singletonList(
+ new UnaryMeasurementSchema(
+ measurementId[0],
+ TSDataType.INT32,
+ TSEncoding.RLE,
+ CompressionType.UNCOMPRESSED,
+ Collections.emptyMap())));
+ List<Pair<Modification, IMemTable>> modsToMemtable = new ArrayList<>();
+ Modification deletion =
+ new Deletion(new PartialPath(deviceId, measurementId[0]),
Long.MAX_VALUE, 10, dataSize);
+ modsToMemtable.add(new Pair<>(deletion, memTable));
+ ReadOnlyMemChunk memChunk = memTable.query(fullPath, Long.MIN_VALUE,
modsToMemtable);
+ IPointReader iterator = memChunk.getPointReader();
+ int cnt = 0;
+ while (iterator.hasNextTimeValuePair()) {
+ TimeValuePair timeValuePair = iterator.nextTimeValuePair();
+ Assert.assertEquals(cnt, timeValuePair.getTimestamp());
+ Assert.assertEquals(cnt,
timeValuePair.getValue().getVector()[0].getInt());
+ cnt++;
+ }
+ Assert.assertEquals(10, cnt);
+ }
+
private void write(
IMemTable memTable,
String deviceId,
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBAggregationWithoutValueFilterWithDeletion2IT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBAggregationWithoutValueFilterWithDeletion2IT.java
index 40f14e4..2ddef8d 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBAggregationWithoutValueFilterWithDeletion2IT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBAggregationWithoutValueFilterWithDeletion2IT.java
@@ -58,9 +58,6 @@ public class IoTDBAggregationWithoutValueFilterWithDeletion2IT
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 21");
} catch (Exception e) {
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBAggregationWithoutValueFilterWithDeletionIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBAggregationWithoutValueFilterWithDeletionIT.java
index bf68c50..cfd647b 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBAggregationWithoutValueFilterWithDeletionIT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBAggregationWithoutValueFilterWithDeletionIT.java
@@ -61,9 +61,6 @@ public class IoTDBAggregationWithoutValueFilterWithDeletionIT
{
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 21");
} catch (Exception e) {
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBDeleteTimeseriesIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBDeleteTimeseriesIT.java
new file mode 100644
index 0000000..b208bc6
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBDeleteTimeseriesIT.java
@@ -0,0 +1,215 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.integration.aligned;
+
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.utils.EnvironmentUtils;
+import org.apache.iotdb.jdbc.Config;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.ResultSetMetaData;
+import java.sql.Statement;
+
+import static org.apache.iotdb.db.constant.TestConstant.TIMESTAMP_STR;
+import static org.apache.iotdb.db.constant.TestConstant.count;
+import static org.junit.Assert.fail;
+
+public class IoTDBDeleteTimeseriesIT {
+
+ private long memtableSizeThreshold;
+
+ @Before
+ public void setUp() throws ClassNotFoundException {
+ Class.forName(Config.JDBC_DRIVER_NAME);
+ EnvironmentUtils.closeStatMonitor();
+ EnvironmentUtils.envSetUp();
+ memtableSizeThreshold =
IoTDBDescriptor.getInstance().getConfig().getMemtableSizeThreshold();
+ IoTDBDescriptor.getInstance().getConfig().setMemtableSizeThreshold(16);
+ }
+
+ @After
+ public void tearDown() throws Exception {
+
IoTDBDescriptor.getInstance().getConfig().setMemtableSizeThreshold(memtableSizeThreshold);
+ EnvironmentUtils.cleanEnv();
+ }
+
+ @Test
+ public void deleteTimeseriesAndCreateDifferentTypeTest() throws Exception {
+ String[] retArray = new String[] {"1,1,", "2,1.1,"};
+ int cnt = 0;
+
+ try (Connection connection =
+ DriverManager.getConnection("jdbc:iotdb://127.0.0.1:6667/",
"root", "root");
+ Statement statement = connection.createStatement()) {
+ statement.execute(
+ "create aligned timeseries root.turbine1.d1(s1 INT64 encoding=PLAIN
compression=SNAPPY, "
+ + "s2 INT64 encoding=PLAIN compression=SNAPPY)");
+ statement.execute("INSERT INTO root.turbine1.d1(timestamp,s1,s2) ALIGNED
VALUES(1,1,2)");
+ boolean hasResult = statement.execute("SELECT s1 FROM root.turbine1.d1");
+ Assert.assertTrue(hasResult);
+ try (ResultSet resultSet = statement.getResultSet()) {
+ ResultSetMetaData resultSetMetaData = resultSet.getMetaData();
+ while (resultSet.next()) {
+ StringBuilder builder = new StringBuilder();
+ for (int i = 1; i <= resultSetMetaData.getColumnCount(); i++) {
+ builder.append(resultSet.getString(i)).append(",");
+ }
+ Assert.assertEquals(retArray[cnt], builder.toString());
+ cnt++;
+ }
+ }
+ statement.execute("DELETE timeseries root.turbine1.d1.s1");
+ statement.execute("INSERT INTO root.turbine1.d1(timestamp,s1) ALIGNED
VALUES(2,1.1)");
+ statement.execute("FLUSH");
+
+ hasResult = statement.execute("SELECT s1 FROM root.turbine1.d1");
+ Assert.assertTrue(hasResult);
+ try (ResultSet resultSet = statement.getResultSet()) {
+ ResultSetMetaData resultSetMetaData = resultSet.getMetaData();
+ while (resultSet.next()) {
+ StringBuilder builder = new StringBuilder();
+ for (int i = 1; i <= resultSetMetaData.getColumnCount(); i++) {
+ builder.append(resultSet.getString(i)).append(",");
+ }
+ Assert.assertEquals(retArray[cnt], builder.toString());
+ cnt++;
+ }
+ }
+ }
+
+ EnvironmentUtils.restartDaemon();
+
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+ boolean hasResult = statement.execute("SELECT * FROM root.**");
+ Assert.assertTrue(hasResult);
+ }
+ }
+
+ @Test
+ public void deleteTimeseriesAndCreateSameTypeTest() throws Exception {
+ String[] retArray = new String[] {"1,1.0,", "2,5.0,"};
+ int cnt = 0;
+
+ try (Connection connection =
+ DriverManager.getConnection("jdbc:iotdb://127.0.0.1:6667/",
"root", "root");
+ Statement statement = connection.createStatement()) {
+ statement.execute(
+ "create aligned timeseries root.turbine1.d1(s1 FLOAT encoding=PLAIN
compression=SNAPPY, "
+ + "s2 INT64 encoding=PLAIN compression=SNAPPY)");
+ statement.execute("INSERT INTO root.turbine1.d1(timestamp,s1,s2) ALIGNED
VALUES(1,1,2)");
+ boolean hasResult = statement.execute("SELECT s1 FROM root.turbine1.d1");
+ Assert.assertTrue(hasResult);
+ try (ResultSet resultSet = statement.getResultSet()) {
+ ResultSetMetaData resultSetMetaData = resultSet.getMetaData();
+ while (resultSet.next()) {
+ StringBuilder builder = new StringBuilder();
+ for (int i = 1; i <= resultSetMetaData.getColumnCount(); i++) {
+ builder.append(resultSet.getString(i)).append(",");
+ }
+ Assert.assertEquals(retArray[cnt], builder.toString());
+ cnt++;
+ }
+ }
+ statement.execute("DELETE timeseries root.turbine1.d1.s1");
+ statement.execute("INSERT INTO root.turbine1.d1(timestamp,s1) ALIGNED
VALUES(2,5)");
+ statement.execute("FLUSH");
+
+ hasResult = statement.execute("SELECT s1 FROM root.turbine1.d1");
+ Assert.assertTrue(hasResult);
+ try (ResultSet resultSet = statement.getResultSet()) {
+ ResultSetMetaData resultSetMetaData = resultSet.getMetaData();
+ while (resultSet.next()) {
+ StringBuilder builder = new StringBuilder();
+ for (int i = 1; i <= resultSetMetaData.getColumnCount(); i++) {
+ builder.append(resultSet.getString(i)).append(",");
+ }
+ Assert.assertEquals(retArray[cnt], builder.toString());
+ cnt++;
+ }
+ }
+ }
+
+ EnvironmentUtils.restartDaemon();
+
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+ boolean hasResult = statement.execute("SELECT * FROM root.**");
+ Assert.assertTrue(hasResult);
+ }
+ }
+
+ @Test
+ public void deleteTimeSeriesMultiIntervalTest() {
+ String[] retArray1 = new String[] {"0,0"};
+
+ int preAvgSeriesPointNumberThreshold =
+
IoTDBDescriptor.getInstance().getConfig().getAvgSeriesPointNumberThreshold();
+ try (Connection connection =
+ DriverManager.getConnection("jdbc:iotdb://127.0.0.1:6667/",
"root", "root");
+ Statement statement = connection.createStatement()) {
+
+
IoTDBDescriptor.getInstance().getConfig().setAvgSeriesPointNumberThreshold(2);
+ String insertSql = "insert into root.sg.d1(time, s1) aligned values(%d,
%d)";
+ for (int i = 1; i <= 4; i++) {
+ statement.execute(String.format(insertSql, i, i));
+ }
+ statement.execute("flush");
+
+ statement.execute("delete from root.sg.d1.s1 where time >= 1 and time <=
2");
+ statement.execute("delete from root.sg.d1.s1 where time >= 3 and time <=
4");
+
+ boolean hasResultSet =
+ statement.execute("select count(s1) from root.sg.d1 where time >= 3
and time <= 4");
+
+ Assert.assertTrue(hasResultSet);
+ int cnt = 0;
+ try (ResultSet resultSet = statement.getResultSet()) {
+ while (resultSet.next()) {
+ String ans =
+ resultSet.getString(TIMESTAMP_STR)
+ + ","
+ + resultSet.getString(count("root.sg.d1.s1"));
+ Assert.assertEquals(retArray1[cnt], ans);
+ cnt++;
+ }
+ Assert.assertEquals(retArray1.length, cnt);
+ }
+ } catch (Exception e) {
+ e.printStackTrace();
+ fail(e.getMessage());
+ } finally {
+ IoTDBDescriptor.getInstance()
+ .getConfig()
+ .setAvgSeriesPointNumberThreshold(preAvgSeriesPointNumberThreshold);
+ }
+ }
+}
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBDeletionIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBDeletionIT.java
new file mode 100644
index 0000000..d296454
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBDeletionIT.java
@@ -0,0 +1,519 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.integration.aligned;
+
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.utils.EnvironmentUtils;
+import org.apache.iotdb.jdbc.Config;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Ignore;
+import org.junit.Test;
+
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.util.Locale;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.fail;
+
+public class IoTDBDeletionIT {
+
+ private static String[] creationSqls =
+ new String[] {
+ "SET STORAGE GROUP TO root.vehicle",
+ "CREATE ALIGNED TIMESERIES root.vehicle.d0(s0 INT32 ENCODING=RLE, s1
INT64 ENCODING=RLE, s2 FLOAT ENCODING=RLE, s3 TEXT ENCODING=PLAIN, s4 BOOLEAN
ENCODING=PLAIN)",
+ };
+
+ private String insertTemplate =
+ "INSERT INTO root.vehicle.d0(timestamp,s0,s1,s2,s3,s4"
+ + ") ALIGNED VALUES(%d,%d,%d,%f,%s,%b)";
+ private String deleteAllTemplate = "DELETE FROM root.vehicle.d0 WHERE time
<= 10000";
+ private long prevPartitionInterval;
+
+ @Before
+ public void setUp() throws Exception {
+ Locale.setDefault(Locale.ENGLISH);
+ EnvironmentUtils.closeStatMonitor();
+ prevPartitionInterval =
IoTDBDescriptor.getInstance().getConfig().getPartitionInterval();
+ IoTDBDescriptor.getInstance().getConfig().setPartitionInterval(1000);
+ EnvironmentUtils.envSetUp();
+ Class.forName(Config.JDBC_DRIVER_NAME);
+ prepareSeries();
+ }
+
+ @After
+ public void tearDown() throws Exception {
+ EnvironmentUtils.cleanEnv();
+
IoTDBDescriptor.getInstance().getConfig().setPartitionInterval(prevPartitionInterval);
+ }
+
+ /**
+ * Should delete this case after the deletion value filter feature be
implemented
+ *
+ * @throws SQLException
+ */
+ @Test
+ public void testUnsupportedValueFilter() throws SQLException {
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ statement.execute("insert into root.vehicle.d0(time,s0) aligned values
(10,310)");
+ statement.execute("insert into root.vehicle.d0(time,s3) aligned values
(10,'text')");
+ statement.execute("insert into root.vehicle.d0(time,s4) aligned values
(10,true)");
+
+ String errorMsg =
+ "303: Check metadata error: For delete statement, where clause can
only"
+ + " contain time expressions, value filter is not currently
supported.";
+
+ String errorMsg2 =
+ "303: Check metadata error: For delete statement, where clause can
only contain"
+ + " atomic expressions like : time > XXX, time <= XXX,"
+ + " or two atomic expressions connected by 'AND'";
+
+ try {
+ statement.execute(
+ "DELETE FROM root.vehicle.d0.s0 WHERE s0 <= 300 AND time > 0 AND
time < 100");
+ fail("should not reach here!");
+ } catch (SQLException e) {
+ assertEquals(errorMsg2, e.getMessage());
+ }
+
+ try {
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE s0 <= 300 AND
s0 > 0");
+ fail("should not reach here!");
+ } catch (SQLException e) {
+ assertEquals(errorMsg, e.getMessage());
+ }
+
+ try {
+ statement.execute("DELETE FROM root.vehicle.d0.s3 WHERE s3 = 'text'");
+ fail("should not reach here!");
+ } catch (SQLException e) {
+ assertEquals(errorMsg, e.getMessage());
+ }
+
+ try {
+ statement.execute("DELETE FROM root.vehicle.d0.s4 WHERE s4 != true");
+ fail("should not reach here!");
+ } catch (SQLException e) {
+ assertEquals(errorMsg, e.getMessage());
+ }
+
+ try (ResultSet set = statement.executeQuery("SELECT s0 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(1, cnt);
+ }
+
+ try (ResultSet set = statement.executeQuery("SELECT s3 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(1, cnt);
+ }
+
+ try (ResultSet set = statement.executeQuery("SELECT s4 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(1, cnt);
+ }
+ }
+ }
+
+ @Test
+ public void test() throws SQLException {
+ prepareData();
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE time <= 300");
+ statement.execute(
+ "DELETE FROM
root.vehicle.d0.s1,root.vehicle.d0.s2,root.vehicle.d0.s3"
+ + " WHERE time <= 350");
+ statement.execute("DELETE FROM root.vehicle.d0.** WHERE time <= 150");
+
+ try (ResultSet set = statement.executeQuery("SELECT * FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(250, cnt);
+ }
+
+ try (ResultSet set = statement.executeQuery("SELECT s0 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(100, cnt);
+ }
+
+ try (ResultSet set = statement.executeQuery("SELECT s1,s2,s3 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(50, cnt);
+ }
+ }
+ cleanData();
+ }
+
+ @Test
+ @Ignore // TODO
+ public void testMerge() throws SQLException {
+ prepareMerge();
+
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+ statement.execute("merge");
+ statement.execute("DELETE FROM root.vehicle.d0.** WHERE time <= 15000");
+
+ // before merge completes
+ try (ResultSet set = statement.executeQuery("SELECT * FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(5000, cnt);
+ }
+
+ // after merge completes
+ try (ResultSet set = statement.executeQuery("SELECT * FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(5000, cnt);
+ }
+ cleanData();
+ }
+ }
+
+ @Test
+ public void testDelAfterFlush() throws SQLException {
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+ statement.execute("SET STORAGE GROUP TO root.ln.wf01.wt01");
+ statement.execute(
+ "CREATE TIMESERIES root.ln.wf01.wt01.status WITH DATATYPE=BOOLEAN,"
+ " ENCODING=PLAIN");
+ statement.execute(
+ "INSERT INTO root.ln.wf01.wt01(timestamp,status) " +
"values(1509465600000,true)");
+ statement.execute("INSERT INTO root.ln.wf01.wt01(timestamp,status)
VALUES(NOW(), false)");
+
+ statement.execute("delete from root.ln.wf01.wt01.status where time <=
NOW()");
+ statement.execute("flush");
+ statement.execute("delete from root.ln.wf01.wt01.status where time <=
NOW()");
+
+ try (ResultSet resultSet = statement.executeQuery("select status from
root.ln.wf01.wt01")) {
+ assertFalse(resultSet.next());
+ }
+ }
+ }
+
+ @Test
+ public void testRangeDelete() throws SQLException {
+ prepareData();
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE time <= 300");
+ statement.execute("DELETE FROM root.vehicle.d0.s1 WHERE time > 150");
+ try (ResultSet set = statement.executeQuery("SELECT s0 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(100, cnt);
+ }
+
+ try (ResultSet set = statement.executeQuery("SELECT s1 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(150, cnt);
+ }
+
+ statement.execute("DELETE FROM root.vehicle.d0.** WHERE time > 50 and
time <= 250");
+ try (ResultSet set = statement.executeQuery("SELECT * FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(200, cnt);
+ }
+ }
+ cleanData();
+ }
+
+ @Test
+ public void testFullDeleteWithoutWhereClause() throws SQLException {
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+ statement.execute("DELETE FROM root.vehicle.d0.s0");
+ try (ResultSet set = statement.executeQuery("SELECT s0 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(0, cnt);
+ }
+ cleanData();
+ }
+ }
+
+ @Test
+ public void testPartialPathRangeDelete() throws SQLException {
+ prepareData();
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ statement.execute("DELETE FROM root.vehicle.d0.* WHERE time <= 300 and
time > 150");
+ try (ResultSet set = statement.executeQuery("SELECT s0 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(250, cnt);
+ }
+
+ statement.execute("DELETE FROM root.vehicle.*.s0 WHERE time <= 100");
+ try (ResultSet set = statement.executeQuery("SELECT s0 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(150, cnt);
+ }
+ }
+ cleanData();
+ }
+
+ @Test
+ public void testDelFlushingMemtable() throws SQLException {
+ long size =
IoTDBDescriptor.getInstance().getConfig().getMemtableSizeThreshold();
+ // Adjust memstable threshold size to make it flush automatically
+ IoTDBDescriptor.getInstance().getConfig().setMemtableSizeThreshold(10000);
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ for (int i = 1; i <= 10000; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE time > 1500 and
time <= 9000");
+ try (ResultSet set = statement.executeQuery("SELECT s0 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(2500, cnt);
+ }
+ cleanData();
+ }
+ IoTDBDescriptor.getInstance().getConfig().setMemtableSizeThreshold(size);
+ }
+
+ @Test
+ public void testDelMultipleFlushingMemtable() throws SQLException {
+ long size =
IoTDBDescriptor.getInstance().getConfig().getMemtableSizeThreshold();
+ // Adjust memstable threshold size to make it flush automatically
+
IoTDBDescriptor.getInstance().getConfig().setMemtableSizeThreshold(1000000);
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ for (int i = 1; i <= 100000; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE time > 15000 and
time <= 30000");
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE time > 30000 and
time <= 40000");
+ for (int i = 100001; i <= 200000; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE time > 50000 and
time <= 80000");
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE time > 90000 and
time <= 110000");
+ statement.execute("DELETE FROM root.vehicle.d0.s0 WHERE time > 150000
and time <= 165000");
+ statement.execute("flush");
+ try (ResultSet set = statement.executeQuery("SELECT s0 FROM
root.vehicle.d0")) {
+ int cnt = 0;
+ while (set.next()) {
+ cnt++;
+ }
+ assertEquals(110000, cnt);
+ }
+ cleanData();
+ }
+ IoTDBDescriptor.getInstance().getConfig().setMemtableSizeThreshold(size);
+ }
+
+ @Test
+ public void testDelSeriesWithSpecialSymbol() throws SQLException {
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+ statement.execute(
+ "CREATE TIMESERIES root.ln.d1.`\"status,01\"` WITH DATATYPE=BOOLEAN,
ENCODING=PLAIN");
+ statement.execute("INSERT INTO root.ln.d1(timestamp,`\"status,01\"`)
VALUES(300, true)");
+ statement.execute("INSERT INTO root.ln.d1(timestamp,`\"status,01\"`)
VALUES(500, false)");
+
+ try (ResultSet resultSet = statement.executeQuery("select
`\"status,01\"` from root.ln.d1")) {
+ int cnt = 0;
+ while (resultSet.next()) {
+ cnt++;
+ }
+ Assert.assertEquals(2, cnt);
+ }
+
+ statement.execute("DELETE FROM root.ln.d1.`\"status,01\"` WHERE time <=
400");
+
+ try (ResultSet resultSet = statement.executeQuery("select
`\"status,01\"` from root.ln.d1")) {
+ int cnt = 0;
+ while (resultSet.next()) {
+ cnt++;
+ }
+ Assert.assertEquals(1, cnt);
+ }
+
+ statement.execute("DELETE FROM root.ln.d1.`\"status,01\"`");
+
+ try (ResultSet resultSet = statement.executeQuery("select
`\"status,01\"` from root.ln.d1")) {
+ int cnt = 0;
+ while (resultSet.next()) {
+ cnt++;
+ }
+ Assert.assertEquals(0, cnt);
+ }
+ }
+ }
+
+ private static void prepareSeries() {
+ String sq = null;
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ for (String sql : creationSqls) {
+ sq = sql;
+ statement.execute(sql);
+ }
+ } catch (Exception e) {
+ System.out.println(sq);
+ e.printStackTrace();
+ }
+ }
+
+ private void prepareData() throws SQLException {
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ // prepare BufferWrite file
+ for (int i = 201; i <= 300; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+ // TODO: merge
+ // statement.execute("merge");
+ // prepare Unseq-File
+ for (int i = 1; i <= 100; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+ // TODO: merge
+ // statement.execute("merge");
+ // prepare BufferWrite cache
+ for (int i = 301; i <= 400; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+ // prepare Overflow cache
+ for (int i = 101; i <= 200; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+ }
+ }
+
+ private void cleanData() throws SQLException {
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+ statement.execute(deleteAllTemplate);
+ }
+ }
+
+ public void prepareMerge() throws SQLException {
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ // prepare BufferWrite data
+ for (int i = 10001; i <= 20000; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+ // prepare Overflow data
+ for (int i = 1; i <= 10000; i++) {
+ statement.execute(
+ String.format(insertTemplate, i, i, i, (double) i, "'" + i + "'",
i % 2 == 0));
+ }
+ }
+ }
+}
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithDeletion2IT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithDeletion2IT.java
index 0f0f5af..e2bb534 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithDeletion2IT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithDeletion2IT.java
@@ -56,9 +56,6 @@ public class IoTDBLastQueryWithDeletion2IT extends
IoTDBLastQueryWithDeletionIT
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 27");
} catch (Exception e) {
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithDeletionIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithDeletionIT.java
index 1ec3eb4..51748ac 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithDeletionIT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithDeletionIT.java
@@ -71,9 +71,6 @@ public class IoTDBLastQueryWithDeletionIT {
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 27");
} catch (Exception e) {
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithoutLastCacheWithDeletion2IT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithoutLastCacheWithDeletion2IT.java
index 5355e57..a2a7edf 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithoutLastCacheWithDeletion2IT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithoutLastCacheWithDeletion2IT.java
@@ -60,9 +60,6 @@ public class IoTDBLastQueryWithoutLastCacheWithDeletion2IT
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 27");
} catch (Exception e) {
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithoutLastCacheWithDeletionIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithoutLastCacheWithDeletionIT.java
index 95b1498..8a4564f 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithoutLastCacheWithDeletionIT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBLastQueryWithoutLastCacheWithDeletionIT.java
@@ -74,9 +74,6 @@ public class IoTDBLastQueryWithoutLastCacheWithDeletionIT {
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 27");
} catch (Exception e) {
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithValueFilterWithDeletion2IT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithValueFilterWithDeletion2IT.java
index 490643a..0726f59 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithValueFilterWithDeletion2IT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithValueFilterWithDeletion2IT.java
@@ -57,11 +57,9 @@ public class IoTDBRawQueryWithValueFilterWithDeletion2IT
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 21");
+ statement.execute("delete from root.sg1.d1.s5 where time <= 31 and time
> 20");
} catch (Exception e) {
e.printStackTrace();
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithValueFilterWithDeletionIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithValueFilterWithDeletionIT.java
index a1ce64c..8d0ca7f 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithValueFilterWithDeletionIT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithValueFilterWithDeletionIT.java
@@ -66,11 +66,9 @@ public class IoTDBRawQueryWithValueFilterWithDeletionIT {
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 21");
+ statement.execute("delete from root.sg1.d1.s5 where time <= 31 and time
> 20");
} catch (Exception e) {
e.printStackTrace();
}
@@ -510,7 +508,6 @@ public class IoTDBRawQueryWithValueFilterWithDeletionIT {
"28,null,28,false,null",
"29,null,29,false,null",
"30,null,30,false,null",
- "31,null,null,null,aligned_test31",
"32,null,null,null,aligned_test32",
"33,null,null,null,aligned_test33",
"36,null,null,null,aligned_test36",
@@ -724,7 +721,6 @@ public class IoTDBRawQueryWithValueFilterWithDeletionIT {
"23,null,false,null,null,true,230000.0",
"24,null,true,null,null,true,null",
"25,null,true,null,null,true,null",
- "31,non_aligned_test31,null,null,aligned_test31,null,null",
"32,non_aligned_test32,null,null,aligned_test32,null,null",
"33,non_aligned_test33,null,null,aligned_test33,null,null",
"34,non_aligned_test34,null,null,aligned_test34,null,null",
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithoutValueFilterWithDeletion2IT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithoutValueFilterWithDeletion2IT.java
index f3377d3..71e0a98 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithoutValueFilterWithDeletion2IT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithoutValueFilterWithDeletion2IT.java
@@ -58,11 +58,9 @@ public class IoTDBRawQueryWithoutValueFilterWithDeletion2IT
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 21");
+ statement.execute("delete from root.sg1.d1.s5 where time <= 31 and time
> 20");
} catch (Exception e) {
e.printStackTrace();
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithoutValueFilterWithDeletionIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithoutValueFilterWithDeletionIT.java
index efa1b4c..5396917 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithoutValueFilterWithDeletionIT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/aligned/IoTDBRawQueryWithoutValueFilterWithDeletionIT.java
@@ -66,11 +66,9 @@ public class IoTDBRawQueryWithoutValueFilterWithDeletionIT {
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
Statement statement = connection.createStatement()) {
- // TODO currently aligned data in memory doesn't support deletion, so we
flush all data to
- // disk before doing deletion
- statement.execute("flush");
statement.execute("delete timeseries root.sg1.d1.s2");
statement.execute("delete from root.sg1.d1.s1 where time <= 21");
+ statement.execute("delete from root.sg1.d1.s5 where time <= 31 and time
> 20");
} catch (Exception e) {
e.printStackTrace();
}
@@ -123,7 +121,6 @@ public class IoTDBRawQueryWithoutValueFilterWithDeletionIT {
"28,null,28,false,null",
"29,null,29,false,null",
"30,null,30,false,null",
- "31,null,null,null,aligned_test31",
"32,null,null,null,aligned_test32",
"33,null,null,null,aligned_test33",
"34,null,null,null,aligned_test34",
@@ -208,7 +205,7 @@ public class IoTDBRawQueryWithoutValueFilterWithDeletionIT {
"28,null,28,false,null,null,null,28,false,null",
"29,null,29,false,null,null,null,29,false,null",
"30,null,30,false,null,null,null,30,false,null",
-
"31,null,null,null,aligned_test31,null,31,null,null,non_aligned_test31",
+ "31,null,null,null,null,null,31,null,null,non_aligned_test31",
"32,null,null,null,aligned_test32,null,32,null,null,non_aligned_test32",
"33,null,null,null,aligned_test33,null,33,null,null,non_aligned_test33",
"34,null,null,null,aligned_test34,null,34,null,null,non_aligned_test34",
@@ -295,7 +292,6 @@ public class IoTDBRawQueryWithoutValueFilterWithDeletionIT {
"28,null,28,false,null",
"29,null,29,false,null",
"30,null,30,false,null",
- "31,null,null,null,aligned_test31",
"32,null,null,null,aligned_test32",
"33,null,null,null,aligned_test33",
};
@@ -365,7 +361,6 @@ public class IoTDBRawQueryWithoutValueFilterWithDeletionIT {
"28,null,false,null",
"29,null,false,null",
"30,null,false,null",
- "31,null,null,aligned_test31",
"32,null,null,aligned_test32",
"33,null,null,aligned_test33",
"34,null,null,aligned_test34",
@@ -494,7 +489,6 @@ public class IoTDBRawQueryWithoutValueFilterWithDeletionIT {
"28,null,false,null",
"29,null,false,null",
"30,null,false,null",
- "31,null,null,aligned_test31",
"32,null,null,aligned_test32",
"33,null,null,aligned_test33",
"34,null,null,aligned_test34",
@@ -559,7 +553,7 @@ public class IoTDBRawQueryWithoutValueFilterWithDeletionIT {
"28,null,false,null,null,false,null",
"29,null,false,null,null,false,null",
"30,null,false,null,null,false,null",
- "31,non_aligned_test31,null,null,aligned_test31,null,null",
+ "31,non_aligned_test31,null,null,null,null,null",
"32,non_aligned_test32,null,null,aligned_test32,null,null",
"33,non_aligned_test33,null,null,aligned_test33,null,null",
"34,non_aligned_test34,null,null,aligned_test34,null,null",