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

Reply via email to