This is an automated email from the ASF dual-hosted git repository.
hxd 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 df22109 [IOTDB-1011] Memtable sort in query (#2144)
df22109 is described below
commit df22109ecbdf3a9bc694d69a7afd46bfaacfa94f
Author: SilverNarcissus <[email protected]>
AuthorDate: Thu Dec 3 01:48:17 2020 +0800
[IOTDB-1011] Memtable sort in query (#2144)
* If the tvlist in a memory Chunk is sorted already, then do not need to
copy the tvlist for queries.
---
.../iotdb/db/engine/flush/MemTableFlushTask.java | 2 +-
.../iotdb/db/engine/memtable/AbstractMemTable.java | 52 ++-
.../db/engine/memtable/IWritableMemChunk.java | 24 +-
.../iotdb/db/engine/memtable/WritableMemChunk.java | 35 +-
.../db/engine/querycontext/ReadOnlyMemChunk.java | 21 +-
.../iotdb/db/utils/datastructure/TVList.java | 128 +++---
.../db/engine/memtable/PrimitiveMemTableTest.java | 2 +-
.../db/integration/IoTDBInsertWithQueryIT.java | 503 +++++++++++++++++++++
8 files changed, 670 insertions(+), 97 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
b/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
index e9f4c92..98a79ce 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
@@ -89,7 +89,7 @@ public class MemTableFlushTask {
long startTime = System.currentTimeMillis();
IWritableMemChunk series =
memTable.getMemTableMap().get(deviceId).get(measurementId);
MeasurementSchema desc = series.getSchema();
- TVList tvList = series.getSortedTVList();
+ TVList tvList = series.getSortedTVListForFlush();
sortTime += System.currentTimeMillis() - startTime;
encodingTaskQueue.add(new Pair<>(tvList, desc));
}
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 5960d43..0ddc384 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
@@ -46,31 +46,23 @@ import
org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
public abstract class AbstractMemTable implements IMemTable {
private final Map<String, Map<String, IWritableMemChunk>> memTableMap;
-
+ /**
+ * The initial value is true because we want calculate the text data size
when recover memTable!!
+ */
+ protected boolean disableMemControl = true;
private long version = Long.MAX_VALUE;
-
private List<Modification> modifications = new ArrayList<>();
-
private int avgSeriesPointNumThreshold =
IoTDBDescriptor.getInstance().getConfig()
.getAvgSeriesPointNumberThreshold();
-
/**
* memory size of data points, including TEXT values
*/
private long memSize = 0;
-
/**
* memory usage of all TVLists memory usage regardless of whether these
TVLists are full,
* including TEXT values
*/
private long tvListRamCost = 0;
-
- /**
- * The initial value is true because we want calculate the text data size
when recover
- * memTable!!
- */
- protected boolean disableMemControl = true;
-
private int seriesNumber = 0;
private long totalPointsNum = 0;
@@ -129,14 +121,16 @@ public abstract class AbstractMemTable implements
IMemTable {
}
Object value = insertRowPlan.getValues()[i];
- memSize +=
MemUtils.getRecordSize(insertRowPlan.getMeasurementMNodes()[i].getSchema().getType(),
value,
- disableMemControl);
+ memSize += MemUtils
+
.getRecordSize(insertRowPlan.getMeasurementMNodes()[i].getSchema().getType(),
value,
+ disableMemControl);
write(insertRowPlan.getDeviceId().getFullPath(),
insertRowPlan.getMeasurements()[i],
insertRowPlan.getMeasurementMNodes()[i].getSchema(),
insertRowPlan.getTime(), value);
}
- totalPointsNum += insertRowPlan.getMeasurements().length -
insertRowPlan.getFailedMeasurementNumber();
+ totalPointsNum +=
+ insertRowPlan.getMeasurements().length -
insertRowPlan.getFailedMeasurementNumber();
}
@Override
@@ -146,8 +140,9 @@ public abstract class AbstractMemTable implements IMemTable
{
try {
write(insertTabletPlan, start, end);
memSize += MemUtils.getRecordSize(insertTabletPlan, start, end,
disableMemControl);
- totalPointsNum += (insertTabletPlan.getMeasurements().length -
insertTabletPlan.getFailedMeasurementNumber())
- * (end - start);
+ totalPointsNum += (insertTabletPlan.getMeasurements().length -
insertTabletPlan
+ .getFailedMeasurementNumber())
+ * (end - start);
} catch (RuntimeException e) {
throw new WriteProcessException(e);
}
@@ -168,8 +163,10 @@ public abstract class AbstractMemTable implements
IMemTable {
if (insertTabletPlan.getColumns()[i] == null) {
continue;
}
- IWritableMemChunk memSeries =
createIfNotExistAndGet(insertTabletPlan.getDeviceId().getFullPath(),
- insertTabletPlan.getMeasurements()[i],
insertTabletPlan.getMeasurementMNodes()[i].getSchema());
+ IWritableMemChunk memSeries = createIfNotExistAndGet(
+ insertTabletPlan.getDeviceId().getFullPath(),
+ insertTabletPlan.getMeasurements()[i],
+ insertTabletPlan.getMeasurementMNodes()[i].getSchema());
memSeries.write(insertTabletPlan.getTimes(),
insertTabletPlan.getColumns()[i],
insertTabletPlan.getDataTypes()[i], start, end);
}
@@ -248,11 +245,14 @@ public abstract class AbstractMemTable implements
IMemTable {
return null;
}
List<TimeRange> deletionList = constructDeletionList(deviceId,
measurement, timeLowerBound);
+
IWritableMemChunk memChunk = memTableMap.get(deviceId).get(measurement);
- TVList chunkCopy = memChunk.getTVList().clone();
+ // get sorted tv list is synchronized so different query can get right
sorted list reference
+ TVList chunkCopy = memChunk.getSortedTVListForQuery();
+ int curSize = chunkCopy.size();
- chunkCopy.setDeletionList(deletionList);
- return new ReadOnlyMemChunk(measurement, dataType, encoding, chunkCopy,
props, getVersion());
+ return new ReadOnlyMemChunk(measurement, dataType, encoding, chunkCopy,
props, getVersion(),
+ curSize, deletionList);
}
private List<TimeRange> constructDeletionList(String deviceId, String
measurement,
@@ -273,7 +273,8 @@ public abstract class AbstractMemTable implements IMemTable
{
}
@Override
- public void delete(PartialPath originalPath, PartialPath devicePath, long
startTimestamp, long endTimestamp) {
+ public void delete(PartialPath originalPath, PartialPath devicePath, long
startTimestamp,
+ long endTimestamp) {
Map<String, IWritableMemChunk> deviceMap =
memTableMap.get(devicePath.getFullPath());
if (deviceMap == null) {
return;
@@ -326,7 +327,10 @@ public abstract class AbstractMemTable implements
IMemTable {
public void release() {
for (Entry<String, Map<String, IWritableMemChunk>> entry :
memTableMap.entrySet()) {
for (Entry<String, IWritableMemChunk> subEntry :
entry.getValue().entrySet()) {
- TVListAllocator.getInstance().release(subEntry.getValue().getTVList());
+ TVList list = subEntry.getValue().getTVList();
+ if (list.getReferenceCount() == 0) {
+ TVListAllocator.getInstance().release(list);
+ }
}
}
}
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 9dc19fd..6fe9844 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
@@ -63,12 +63,28 @@ public interface IWritableMemChunk {
/**
* served for query requests.
+ * <p>
+ * if tv list has been sorted, just return reference of it
+ * <p>
+ * if tv list hasn't been sorted and has no reference, sort and return
reference of it
+ * <p>
+ * if tv list hasn't been sorted and has reference we should copy and sort
it, then return ths
+ * list
+ * <p>
+ * the mechanism is just like copy on write
*
- * @return
+ * @return sorted tv list
*/
- default TVList getSortedTVList() {
- return null;
- }
+ TVList getSortedTVListForQuery();
+
+
+ /**
+ * served for flush requests.
+ * The logic is just same as getSortedTVListForQuery, but without add
reference count
+ *
+ * @return sorted tv list
+ */
+ TVList getSortedTVListForFlush();
default TVList getTVList() {
return null;
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 b9f885b..855d567 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
@@ -153,8 +153,33 @@ public class WritableMemChunk implements IWritableMemChunk
{
}
@Override
- public synchronized TVList getSortedTVList() {
- list.sort();
+ public synchronized TVList getSortedTVListForQuery() {
+ // check reference count
+ if ((list.getReferenceCount() > 0 && !list.isSorted())) {
+ list = list.clone();
+ }
+
+ if (!list.isSorted()) {
+ list.sort();
+ }
+
+ // increase reference count
+ list.increaseReferenceCount();
+
+ return list;
+ }
+
+ @Override
+ public synchronized TVList getSortedTVListForFlush() {
+ // check reference count
+ if ((list.getReferenceCount() > 0 && !list.isSorted())) {
+ list = list.clone();
+ }
+
+ if (!list.isSorted()) {
+ list.sort();
+ }
+
return list;
}
@@ -185,13 +210,13 @@ public class WritableMemChunk implements
IWritableMemChunk {
@Override
public String toString() {
- int size = getSortedTVList().size();
+ int size = getSortedTVListForQuery().size();
StringBuilder out = new StringBuilder("MemChunk Size: " + size +
System.lineSeparator());
if (size != 0) {
out.append("Data
type:").append(schema.getType()).append(System.lineSeparator());
- out.append("First point:").append(getSortedTVList().getTimeValuePair(0))
+ out.append("First
point:").append(getSortedTVListForQuery().getTimeValuePair(0))
.append(System.lineSeparator());
- out.append("Last point:").append(getSortedTVList().getTimeValuePair(size
- 1))
+ out.append("Last
point:").append(getSortedTVListForQuery().getTimeValuePair(size - 1))
.append(System.lineSeparator());
}
return out.toString();
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/querycontext/ReadOnlyMemChunk.java
b/server/src/main/java/org/apache/iotdb/db/engine/querycontext/ReadOnlyMemChunk.java
index 8890923..d87a4ea 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/querycontext/ReadOnlyMemChunk.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/querycontext/ReadOnlyMemChunk.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.engine.querycontext;
import java.io.IOException;
+import java.util.List;
import java.util.Map;
import org.apache.iotdb.db.exception.query.QueryProcessException;
import org.apache.iotdb.db.query.reader.chunk.MemChunkLoader;
@@ -30,12 +31,16 @@ import
org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics;
import org.apache.iotdb.tsfile.read.TimeValuePair;
+import org.apache.iotdb.tsfile.read.common.TimeRange;
import org.apache.iotdb.tsfile.read.reader.IPointReader;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class ReadOnlyMemChunk {
+ // deletion list for this chunk
+ private final List<TimeRange> deletionList;
+
private String measurementUid;
private TSDataType dataType;
private TSEncoding encoding;
@@ -51,8 +56,11 @@ public class ReadOnlyMemChunk {
private IPointReader chunkPointReader;
+ private int chunkDataSize;
+
public ReadOnlyMemChunk(String measurementUid, TSDataType dataType,
TSEncoding encoding,
- TVList tvList, Map<String, String> props, long version)
+ TVList tvList, Map<String, String> props, long version, int size,
+ List<TimeRange> deletionList)
throws IOException, QueryProcessException {
this.measurementUid = measurementUid;
this.dataType = dataType;
@@ -72,9 +80,12 @@ public class ReadOnlyMemChunk {
floatPrecision =
TSFileDescriptor.getInstance().getConfig().getFloatPrecision();
}
}
- tvList.sort();
+
this.chunkData = tvList;
- this.chunkPointReader = tvList.getIterator(floatPrecision, encoding);
+ this.chunkDataSize = size;
+ this.deletionList = deletionList;
+
+ this.chunkPointReader = tvList.getIterator(floatPrecision, encoding,
chunkDataSize, deletionList);
initChunkMeta();
}
@@ -82,7 +93,7 @@ public class ReadOnlyMemChunk {
Statistics statsByType = Statistics.getStatsByType(dataType);
ChunkMetadata metaData = new ChunkMetadata(measurementUid, dataType, 0,
statsByType);
if (!isEmpty()) {
- IPointReader iterator = chunkData.getIterator(floatPrecision, encoding);
+ IPointReader iterator = chunkData.getIterator(floatPrecision, encoding,
chunkDataSize, deletionList);
while (iterator.hasNextTimeValuePair()) {
TimeValuePair timeValuePair = iterator.nextTimeValuePair();
switch (dataType) {
@@ -128,7 +139,7 @@ public class ReadOnlyMemChunk {
}
public IPointReader getPointReader() {
- chunkPointReader = chunkData.getIterator(floatPrecision, encoding);
+ chunkPointReader = chunkData.getIterator(floatPrecision, encoding,
chunkDataSize, deletionList);
return chunkPointReader;
}
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 18e7176..9978529 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
@@ -25,7 +25,9 @@ import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
+import java.util.concurrent.atomic.AtomicInteger;
import org.apache.iotdb.db.rescon.PrimitiveArrayManager;
+import org.apache.iotdb.db.utils.TestOnly;
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;
@@ -35,31 +37,69 @@ import org.apache.iotdb.tsfile.utils.Binary;
public abstract class TVList {
- private static final String ERR_DATATYPE_NOT_CONSISTENT = "DataType not
consistent";
-
protected static final int SMALL_ARRAY_LENGTH = 32;
-
+ private static final String ERR_DATATYPE_NOT_CONSISTENT = "DataType not
consistent";
protected List<long[]> timestamps;
protected int size;
protected long[][] sortedTimestamps;
protected boolean sorted = true;
-
- /**
- * this field is effective only in the Tvlist in a RealOnlyMemChunk.
- */
- private List<TimeRange> deletionList;
- private long version;
-
+ // record reference count of this tv list
+ // currently this reference will only be increase because we can't know when
to decrease it
+ protected AtomicInteger referenceCount;
protected long pivotTime;
-
protected long minTime;
+ private long version;
public TVList() {
timestamps = new ArrayList<>();
size = 0;
minTime = Long.MAX_VALUE;
+ referenceCount = new AtomicInteger();
+ }
+
+ public static TVList newList(TSDataType dataType) {
+ switch (dataType) {
+ case TEXT:
+ return new BinaryTVList();
+ case FLOAT:
+ return new FloatTVList();
+ case INT32:
+ return new IntTVList();
+ case INT64:
+ return new LongTVList();
+ case DOUBLE:
+ return new DoubleTVList();
+ case BOOLEAN:
+ return new BooleanTVList();
+ default:
+ break;
+ }
+ return null;
+ }
+
+ public static long tvListArrayMemSize(TSDataType type) {
+ long size = 0;
+ // time size
+ size +=
+ PrimitiveArrayManager.ARRAY_SIZE * 8;
+ // value size
+ size +=
+ PrimitiveArrayManager.ARRAY_SIZE * type.getDataTypeSize();
+ return size;
+ }
+
+ public boolean isSorted() {
+ return sorted;
+ }
+
+ public void increaseReferenceCount() {
+ referenceCount.incrementAndGet();
+ }
+
+ public int getReferenceCount() {
+ return referenceCount.get();
}
public int size() {
@@ -222,9 +262,6 @@ public abstract class TVList {
clearValue();
clearSortedValue();
- if (deletionList != null) {
- deletionList.clear();
- }
}
protected void clearTime() {
@@ -271,6 +308,7 @@ public abstract class TVList {
if (sorted) {
return;
}
+
if (lo == hi) {
return;
}
@@ -307,33 +345,6 @@ public abstract class TVList {
return runHi - lo;
}
- public static TVList newList(TSDataType dataType) {
- switch (dataType) {
- case TEXT:
- return new BinaryTVList();
- case FLOAT:
- return new FloatTVList();
- case INT32:
- return new IntTVList();
- case INT64:
- return new LongTVList();
- case DOUBLE:
- return new DoubleTVList();
- case BOOLEAN:
- return new BooleanTVList();
- default:
- break;
- }
- return null;
- }
-
- /**
- * this field is effective only in the Tvlist in a RealOnlyMemChunk.
- */
- public void setDeletionList(List<TimeRange> list) {
- this.deletionList = list;
- }
-
protected int compare(int idx1, int idx2) {
long t1 = getTime(idx1);
long t2 = getTime(idx2);
@@ -469,23 +480,14 @@ public abstract class TVList {
protected abstract TimeValuePair getTimeValuePair(int index, long time,
Integer floatPrecision, TSEncoding encoding);
+ @TestOnly
public IPointReader getIterator() {
return new Ite();
}
- public IPointReader getIterator(int floatPrecision, TSEncoding encoding) {
- return new Ite(floatPrecision, encoding);
- }
-
- public static long tvListArrayMemSize(TSDataType type) {
- long size = 0;
- // time size
- size +=
- PrimitiveArrayManager.ARRAY_SIZE * 8;
- // value size
- size +=
- PrimitiveArrayManager.ARRAY_SIZE * type.getDataTypeSize();
- return size;
+ public IPointReader getIterator(int floatPrecision, TSEncoding encoding, int
size,
+ List<TimeRange> deletionList) {
+ return new Ite(floatPrecision, encoding, size, deletionList);
}
private class Ite implements IPointReader {
@@ -496,13 +498,24 @@ public abstract class TVList {
private Integer floatPrecision;
private TSEncoding encoding;
private int deleteCursor = 0;
+ /**
+ * because TV list may be share with different query, each iterator has to
record it's own size
+ */
+ private int iteSize = 0;
+ /**
+ * this field is effective only in the Tvlist in a RealOnlyMemChunk.
+ */
+ private List<TimeRange> deletionList;
public Ite() {
+ this.iteSize = TVList.this.size;
}
- public Ite(int floatPrecision, TSEncoding encoding) {
+ public Ite(int floatPrecision, TSEncoding encoding, int size,
List<TimeRange> deletionList) {
this.floatPrecision = floatPrecision;
this.encoding = encoding;
+ this.iteSize = size;
+ this.deletionList = deletionList;
}
@Override
@@ -511,7 +524,7 @@ public abstract class TVList {
return true;
}
- while (cur < size) {
+ while (cur < iteSize) {
long time = getTime(cur);
if (isPointDeleted(time) || (cur + 1 < size() && (time == getTime(cur
+ 1)))) {
cur++;
@@ -522,7 +535,8 @@ public abstract class TVList {
cur++;
return true;
}
- return hasCachedPair;
+
+ return false;
}
private boolean isPointDeleted(long timestamp) {
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 94c6378..b967fe1 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
@@ -58,7 +58,7 @@ public class PrimitiveMemTableTest {
for (int i = 0; i < count; i++) {
series.write(i, i);
}
- IPointReader it = series.getSortedTVList().getIterator();
+ IPointReader it = series.getSortedTVListForQuery().getIterator();
int i = 0;
while (it.hasNextTimeValuePair()) {
Assert.assertEquals(i, it.nextTimeValuePair().getTimestamp());
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBInsertWithQueryIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBInsertWithQueryIT.java
new file mode 100644
index 0000000..9be2193
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBInsertWithQueryIT.java
@@ -0,0 +1,503 @@
+/*
+ * 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;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.fail;
+
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.iotdb.db.constant.TestConstant;
+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;
+
+/**
+ * Notice that, all test begins with "IoTDB" is integration test. All test
which will start the
+ * IoTDB server should be defined as integration test.
+ */
+public class IoTDBInsertWithQueryIT {
+
+
+ @Before
+ public void setUp() {
+
+ EnvironmentUtils.closeStatMonitor();
+
+ EnvironmentUtils.envSetUp();
+ }
+
+ @After
+ public void tearDown() throws Exception {
+ EnvironmentUtils.cleanEnv();
+ }
+
+
+ @Test
+ public void insertWithQueryTest() throws ClassNotFoundException {
+ // insert
+ insertData(0, 1000);
+
+ // select
+ selectAndCount(1000);
+
+ // insert
+ insertData(1000, 2000);
+
+ // select
+ selectAndCount(2000);
+ }
+
+ @Test
+ public void insertWithQueryMultiThreadTest() throws ClassNotFoundException,
InterruptedException {
+ // insert
+ insertData(0, 1000);
+
+ selectWithMultiThread(1000);
+
+ // insert
+ insertData(1000, 2000);
+
+ // select
+ selectWithMultiThread(2000);
+ }
+
+ @Test
+ public void insertWithQueryUnsequenceTest() throws ClassNotFoundException {
+ // insert
+ insertData(0, 1000);
+
+ // select
+ selectAndCount(1000);
+
+ // insert
+ insertData(500, 1500);
+
+ // select
+ selectAndCount(1500);
+
+ // insert
+ insertData(2000, 3000);
+
+ // select
+ selectAndCount(2500);
+ }
+
+ @Test
+ public void insertWithQueryMultiThreadUnsequenceTest()
+ throws ClassNotFoundException, InterruptedException {
+ // insert
+ insertData(0, 1000);
+
+ selectWithMultiThread(1000);
+
+ // insert
+ insertData(500, 1500);
+
+ // select
+ selectWithMultiThread(1500);
+
+ // insert
+ insertData(2000, 3000);
+
+ // select
+ selectWithMultiThread(2500);
+ }
+
+ @Test
+ public void insertWithQueryFlushTest() throws ClassNotFoundException {
+ // insert
+ insertData(0, 1000);
+
+ // select
+ selectAndCount(1000);
+
+ flush();
+
+ // insert
+ insertData(1000, 2000);
+
+ // select
+ selectAndCount(2000);
+ }
+
+ @Test
+ public void flushWithQueryTest() throws ClassNotFoundException,
InterruptedException {
+ // insert
+ insertData(0, 1000);
+
+ // select with flush
+ selectWithMultiThreadAndFlush(1000);
+
+ // insert
+ insertData(500, 1500);
+
+ // select
+ selectWithMultiThreadAndFlush(1500);
+ }
+
+ @Test
+ public void flushWithQueryUnorderTest() throws ClassNotFoundException,
InterruptedException {
+ // insert
+ insertData(0, 100);
+ insertData(500, 600);
+
+ // select
+ selectWithMultiThread(200);
+
+ insertData(200, 400);
+
+ selectWithMultiThreadAndFlush(400);
+
+ insertData(0, 1000);
+
+ selectWithMultiThread(1000);
+ }
+
+ @Test
+ public void flushWithQueryUnorderLargerTest() throws ClassNotFoundException,
InterruptedException {
+ // insert
+ insertData(0, 100);
+ insertData(500, 600);
+
+ // select
+ selectWithMultiThread(200);
+
+ insertData(200, 400);
+
+ selectWithMultiThreadAndFlush(400);
+
+ insertData(400, 700);
+//
+ selectWithMultiThreadAndFlush(600);
+
+ insertData(0, 1000);
+
+ selectWithMultiThread(1000);
+//
+ insertData(800, 1500);
+
+ selectWithMultiThreadAndFlush(1500);
+ }
+
+ @Test
+ public void insertWithQueryTogetherTest() throws InterruptedException {
+ // insert
+ List<Thread> queryThreadList = new ArrayList<>();
+
+ // select with multi thread
+ Thread cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ insertData(0, 200);
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ insertData(200, 400);
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ select();
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ insertData(100, 200);
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ select();
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ insertData(700, 900);
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ select();
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ flush();
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ insertData(500, 700);
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ select();
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+ queryThreadList.add(cur);
+ cur.start();
+
+ for (Thread thread : queryThreadList) {
+ thread.join();
+ }
+ }
+
+
+ private void selectWithMultiThreadAndFlush(int res) throws
InterruptedException {
+ List<Thread> queryThreadList = new ArrayList<>();
+
+ // select with multi thread
+ for (int i = 0; i < 5; i++) {
+ Thread cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ selectAndCount(res);
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+
+ if(i == 2){
+ Thread flushThread = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ flush();
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+
+ flushThread.start();
+ queryThreadList.add(flushThread);
+ }
+
+ queryThreadList.add(cur);
+ cur.start();
+ }
+
+ for (Thread thread : queryThreadList) {
+ thread.join();
+ }
+ }
+
+
+ private void selectWithMultiThread(int res) throws InterruptedException {
+ List<Thread> queryThreadList = new ArrayList<>();
+
+ // select with multi thread
+ for (int i = 0; i < 5; i++) {
+ Thread cur = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ selectAndCount(res);
+ } catch (ClassNotFoundException e) {
+ e.printStackTrace();
+ }
+ }
+ });
+
+ queryThreadList.add(cur);
+ cur.start();
+ }
+
+ for (Thread thread : queryThreadList) {
+ thread.join();
+ }
+ }
+
+ private void insertData(int start, int end) throws ClassNotFoundException {
+ Class.forName(Config.JDBC_DRIVER_NAME);
+ try (Connection connection = DriverManager
+ .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
"root");
+ Statement statement = connection.createStatement()) {
+ // insert of data time range : start-end into fans
+ for (int time = start; time < end; time++) {
+ String sql = String
+ .format("insert into root.fans.d0(timestamp,s0) values(%s,%s)",
time, time % 70);
+ statement.execute(sql);
+ sql = String
+ .format("insert into root.fans.d0(timestamp,s1) values(%s,%s)",
time, time % 40);
+ statement.execute(sql);
+ }
+ } catch (SQLException e) {
+ e.printStackTrace();
+ }
+ }
+
+ private void flush() throws ClassNotFoundException {
+ Class.forName(Config.JDBC_DRIVER_NAME);
+ try (Connection connection = DriverManager
+ .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
"root");
+ Statement statement = connection.createStatement()) {
+ // insert of data time range : start-end into fans
+ statement.execute("flush");
+ } catch (SQLException e) {
+ e.printStackTrace();
+ }
+ }
+
+ // test count
+ private void selectAndCount(int res) throws ClassNotFoundException {
+ String selectSql = "select * from root";
+
+ Class.forName(Config.JDBC_DRIVER_NAME);
+ try (Connection connection = DriverManager
+ .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
"root");
+ Statement statement = connection.createStatement()) {
+ boolean hasResultSet = statement.execute(selectSql);
+ Assert.assertTrue(hasResultSet);
+ try (ResultSet resultSet = statement.getResultSet()) {
+ int cnt = 0;
+ long before = -1;
+ while (resultSet.next()) {
+ long cur =
Long.parseLong(resultSet.getString(TestConstant.TIMESTAMP_STR));
+ if(cur <= before){
+ fail("time order is wrong");
+ }
+ before = cur;
+ cnt++;
+ }
+ assertEquals(res, cnt);
+ }
+ } catch (Exception e) {
+ e.printStackTrace();
+ fail(e.getMessage());
+ }
+ }
+
+ // test order
+ private void select() throws ClassNotFoundException {
+ String selectSql = "select * from root";
+
+ Class.forName(Config.JDBC_DRIVER_NAME);
+ try (Connection connection = DriverManager
+ .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
"root");
+ Statement statement = connection.createStatement()) {
+ boolean hasResultSet = statement.execute(selectSql);
+ Assert.assertTrue(hasResultSet);
+ try (ResultSet resultSet = statement.getResultSet()) {
+ int cnt = 0;
+ long before = -1;
+ while (resultSet.next()) {
+ long cur =
Long.parseLong(resultSet.getString(TestConstant.TIMESTAMP_STR));
+ if(cur <= before){
+ fail("time order is wrong");
+ }
+ before = cur;
+ cnt++;
+ }
+ }
+ } catch (Exception e) {
+ e.printStackTrace();
+ fail(e.getMessage());
+ }
+ }
+}