This is an automated email from the ASF dual-hosted git repository.

haonan 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 def138b590 [IOTDB-3971] Improve the process of writing chunks in 
compaction for aligned series (#6786)
def138b590 is described below

commit def138b5909e5df76440836719b62ebc4b48daeb
Author: Mrquan <[email protected]>
AuthorDate: Fri Jul 29 18:14:29 2022 +0800

    [IOTDB-3971] Improve the process of writing chunks in compaction for 
aligned series (#6786)
---
 .../impl/ReadPointCompactionPerformer.java         | 18 +++---
 .../writer/AbstractCompactionWriter.java           |  7 ++-
 .../writer/CrossSpaceCompactionWriter.java         | 19 +++++-
 .../writer/InnerSpaceCompactionWriter.java         | 15 ++++-
 .../file/metadata/statistics/Statistics.java       | 10 +++
 .../file/metadata/statistics/TimeStatistics.java   |  8 +++
 .../iotdb/tsfile/read/common/block/TsBlock.java    |  4 ++
 .../read/common/block/column/BinaryColumn.java     | 16 +++++
 .../read/common/block/column/BooleanColumn.java    | 16 +++++
 .../tsfile/read/common/block/column/Column.java    | 38 ++++++++++++
 .../read/common/block/column/DoubleColumn.java     | 16 +++++
 .../read/common/block/column/FloatColumn.java      | 16 +++++
 .../tsfile/read/common/block/column/IntColumn.java | 16 +++++
 .../read/common/block/column/LongColumn.java       | 16 +++++
 .../block/column/RunLengthEncodedColumn.java       | 58 +++++++++++++++++
 .../read/common/block/column/TimeColumn.java       | 11 ++++
 .../tsfile/write/chunk/AlignedChunkWriterImpl.java | 63 +++++++++++++++++++
 .../iotdb/tsfile/write/chunk/TimeChunkWriter.java  | 10 ++-
 .../iotdb/tsfile/write/chunk/ValueChunkWriter.java | 24 ++++----
 .../iotdb/tsfile/write/page/TimePageWriter.java    |  6 +-
 .../iotdb/tsfile/write/page/ValuePageWriter.java   | 72 ++++++++++++++--------
 21 files changed, 404 insertions(+), 55 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/performer/impl/ReadPointCompactionPerformer.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/performer/impl/ReadPointCompactionPerformer.java
index a47574a7b0..81e29aa7c6 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/performer/impl/ReadPointCompactionPerformer.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/performer/impl/ReadPointCompactionPerformer.java
@@ -371,15 +371,19 @@ public class ReadPointCompactionPerformer
       throws IOException {
     while (reader.hasNextBatch()) {
       TsBlock tsBlock = reader.nextBatch();
-      IPointReader pointReader;
       if (isAligned) {
-        pointReader = tsBlock.getTsBlockAlignedRowIterator();
+        writer.write(
+            tsBlock.getTimeColumn(),
+            tsBlock.getValueColumns(),
+            subTaskId,
+            tsBlock.getPositionCount());
       } else {
-        pointReader = tsBlock.getTsBlockSingleColumnIterator();
-      }
-      while (pointReader.hasNextTimeValuePair()) {
-        TimeValuePair timeValuePair = pointReader.nextTimeValuePair();
-        writer.write(timeValuePair.getTimestamp(), 
timeValuePair.getValue().getValue(), subTaskId);
+        IPointReader pointReader = tsBlock.getTsBlockSingleColumnIterator();
+        while (pointReader.hasNextTimeValuePair()) {
+          TimeValuePair timeValuePair = pointReader.nextTimeValuePair();
+          writer.write(
+              timeValuePair.getTimestamp(), 
timeValuePair.getValue().getValue(), subTaskId);
+        }
       }
     }
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/AbstractCompactionWriter.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/AbstractCompactionWriter.java
index 343e88ee58..6da38ebc16 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/AbstractCompactionWriter.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/AbstractCompactionWriter.java
@@ -25,6 +25,8 @@ import 
org.apache.iotdb.db.engine.compaction.constant.ProcessChunkType;
 import org.apache.iotdb.db.service.metrics.recorder.CompactionMetricsRecorder;
 import org.apache.iotdb.metrics.config.MetricConfigDescriptor;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.read.common.block.column.Column;
+import org.apache.iotdb.tsfile.read.common.block.column.TimeColumn;
 import org.apache.iotdb.tsfile.utils.Binary;
 import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 import org.apache.iotdb.tsfile.write.chunk.AlignedChunkWriterImpl;
@@ -73,7 +75,8 @@ public abstract class AbstractCompactionWriter implements 
AutoCloseable {
 
   public abstract void write(long timestamp, Object value, int subTaskId) 
throws IOException;
 
-  public abstract void write(long[] timestamps, Object values);
+  public abstract void write(TimeColumn timestamps, Column[] columns, int 
subTaskId, int batchSize)
+      throws IOException;
 
   public abstract void endFile() throws IOException;
 
@@ -150,7 +153,7 @@ public abstract class AbstractCompactionWriter implements 
AutoCloseable {
 
   protected void checkChunkSizeAndMayOpenANewChunk(TsFileIOWriter fileWriter, 
int subTaskId)
       throws IOException {
-    if (measurementPointCountArray[subTaskId] % 10 == 0 && 
checkChunkSize(subTaskId)) {
+    if (checkChunkSize(subTaskId)) {
       flushChunkToFileWriter(fileWriter, subTaskId);
       CompactionMetricsRecorder.recordWriteInfo(
           this instanceof CrossSpaceCompactionWriter
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/CrossSpaceCompactionWriter.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/CrossSpaceCompactionWriter.java
index 3e245cfc35..80902dd1d9 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/CrossSpaceCompactionWriter.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/CrossSpaceCompactionWriter.java
@@ -21,6 +21,9 @@ package org.apache.iotdb.db.engine.compaction.writer;
 import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
 import org.apache.iotdb.db.query.control.FileReaderManager;
 import org.apache.iotdb.tsfile.file.metadata.TimeseriesMetadata;
+import org.apache.iotdb.tsfile.read.common.block.column.Column;
+import org.apache.iotdb.tsfile.read.common.block.column.TimeColumn;
+import org.apache.iotdb.tsfile.write.chunk.AlignedChunkWriterImpl;
 import org.apache.iotdb.tsfile.write.writer.TsFileIOWriter;
 
 import java.io.IOException;
@@ -99,13 +102,25 @@ public class CrossSpaceCompactionWriter extends 
AbstractCompactionWriter {
   public void write(long timestamp, Object value, int subTaskId) throws 
IOException {
     checkTimeAndMayFlushChunkToCurrentFile(timestamp, subTaskId);
     writeDataPoint(timestamp, value, subTaskId);
-    
checkChunkSizeAndMayOpenANewChunk(fileWriterList.get(seqFileIndexArray[subTaskId]),
 subTaskId);
+    if (measurementPointCountArray[subTaskId] % 10 == 0) {
+      checkChunkSizeAndMayOpenANewChunk(
+          fileWriterList.get(seqFileIndexArray[subTaskId]), subTaskId);
+    }
     isDeviceExistedInTargetFiles[seqFileIndexArray[subTaskId]] = true;
     isEmptyFile[seqFileIndexArray[subTaskId]] = false;
   }
 
   @Override
-  public void write(long[] timestamps, Object values) {}
+  public void write(TimeColumn timestamps, Column[] columns, int subTaskId, 
int batchSize)
+      throws IOException {
+    // todo control time range of target tsfile
+    checkTimeAndMayFlushChunkToCurrentFile(timestamps.getStartTime(), 
subTaskId);
+    AlignedChunkWriterImpl chunkWriter = (AlignedChunkWriterImpl) 
this.chunkWriters[subTaskId];
+    chunkWriter.write(timestamps, columns, batchSize);
+    
checkChunkSizeAndMayOpenANewChunk(fileWriterList.get(seqFileIndexArray[subTaskId]),
 subTaskId);
+    isDeviceExistedInTargetFiles[seqFileIndexArray[subTaskId]] = true;
+    isEmptyFile[seqFileIndexArray[subTaskId]] = false;
+  }
 
   @Override
   public void endFile() throws IOException {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/InnerSpaceCompactionWriter.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/InnerSpaceCompactionWriter.java
index af2cc53c67..a73c6c2907 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/InnerSpaceCompactionWriter.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/writer/InnerSpaceCompactionWriter.java
@@ -19,6 +19,9 @@
 package org.apache.iotdb.db.engine.compaction.writer;
 
 import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.tsfile.read.common.block.column.Column;
+import org.apache.iotdb.tsfile.read.common.block.column.TimeColumn;
+import org.apache.iotdb.tsfile.write.chunk.AlignedChunkWriterImpl;
 import org.apache.iotdb.tsfile.write.writer.TsFileIOWriter;
 
 import java.io.IOException;
@@ -55,12 +58,20 @@ public class InnerSpaceCompactionWriter extends 
AbstractCompactionWriter {
   @Override
   public void write(long timestamp, Object value, int subTaskId) throws 
IOException {
     writeDataPoint(timestamp, value, subTaskId);
-    checkChunkSizeAndMayOpenANewChunk(fileWriter, subTaskId);
+    if (measurementPointCountArray[subTaskId] % 10 == 0) {
+      checkChunkSizeAndMayOpenANewChunk(fileWriter, subTaskId);
+    }
     isEmptyFile = false;
   }
 
   @Override
-  public void write(long[] timestamps, Object values) {}
+  public void write(TimeColumn timestamps, Column[] columns, int subTaskId, 
int batchSize)
+      throws IOException {
+    AlignedChunkWriterImpl chunkWriter = (AlignedChunkWriterImpl) 
this.chunkWriters[subTaskId];
+    chunkWriter.write(timestamps, columns, batchSize);
+    checkChunkSizeAndMayOpenANewChunk(fileWriter, subTaskId);
+    isEmptyFile = false;
+  }
 
   @Override
   public void endFile() throws IOException {
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/Statistics.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/Statistics.java
index 197d9848f8..1ba48d311b 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/Statistics.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/Statistics.java
@@ -257,6 +257,16 @@ public abstract class Statistics<T extends Serializable> {
     count += batchSize;
   }
 
+  public void update(long[] time, int batchSize, int arrayOffset) {
+    if (time[arrayOffset] < startTime) {
+      startTime = time[arrayOffset];
+    }
+    if (time[arrayOffset + batchSize - 1] > this.endTime) {
+      endTime = time[arrayOffset + batchSize - 1];
+    }
+    count += batchSize;
+  }
+
   protected abstract void mergeStatisticsValue(Statistics<T> stats);
 
   public boolean isEmpty() {
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/TimeStatistics.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/TimeStatistics.java
index 33fcad15cb..60ce558187 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/TimeStatistics.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/statistics/TimeStatistics.java
@@ -59,6 +59,14 @@ public class TimeStatistics extends Statistics<Long> {
     }
   }
 
+  @Override
+  public void update(long[] time, int batchSize, int arrayOffset) {
+    super.update(time, batchSize, arrayOffset);
+    if (batchSize > 0) {
+      setEmpty(false);
+    }
+  }
+
   @Override
   public Long getMinValue() {
     throw new StatisticsClassException(String.format(STATS_UNSUPPORTED_MSG, 
TIME, "min value"));
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/TsBlock.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/TsBlock.java
index c7206efb57..20d408abc6 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/TsBlock.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/TsBlock.java
@@ -198,6 +198,10 @@ public class TsBlock {
     return timeColumn;
   }
 
+  public Column[] getValueColumns() {
+    return valueColumns;
+  }
+
   public Column getColumn(int columnIndex) {
     return valueColumns[columnIndex];
   }
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/BinaryColumn.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/BinaryColumn.java
index 80fba8b7d2..21d822da9d 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/BinaryColumn.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/BinaryColumn.java
@@ -24,6 +24,7 @@ import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
 import org.openjdk.jol.info.ClassLayout;
 
+import java.util.Arrays;
 import java.util.Optional;
 
 import static io.airlift.slice.SizeOf.sizeOf;
@@ -85,6 +86,11 @@ public class BinaryColumn implements Column {
     return values[position + arrayOffset];
   }
 
+  @Override
+  public Binary[] getBinaries() {
+    return values;
+  }
+
   @Override
   public Object getObject(int position) {
     return getBinary(position);
@@ -107,6 +113,16 @@ public class BinaryColumn implements Column {
     return valueIsNull != null && valueIsNull[position + arrayOffset];
   }
 
+  @Override
+  public boolean[] isNull() {
+    if (valueIsNull == null) {
+      boolean[] res = new boolean[positionCount];
+      Arrays.fill(res, false);
+      return res;
+    }
+    return valueIsNull;
+  }
+
   @Override
   public int getPositionCount() {
     return positionCount;
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/BooleanColumn.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/BooleanColumn.java
index 4fab293acb..cf2e6a4f2d 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/BooleanColumn.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/BooleanColumn.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
 import org.openjdk.jol.info.ClassLayout;
 
+import java.util.Arrays;
 import java.util.Optional;
 
 import static io.airlift.slice.SizeOf.sizeOf;
@@ -84,6 +85,11 @@ public class BooleanColumn implements Column {
     return values[position + arrayOffset];
   }
 
+  @Override
+  public boolean[] getBooleans() {
+    return values;
+  }
+
   @Override
   public Object getObject(int position) {
     return getBoolean(position);
@@ -106,6 +112,16 @@ public class BooleanColumn implements Column {
     return valueIsNull != null && valueIsNull[position + arrayOffset];
   }
 
+  @Override
+  public boolean[] isNull() {
+    if (valueIsNull == null) {
+      boolean[] res = new boolean[positionCount];
+      Arrays.fill(res, false);
+      return res;
+    }
+    return valueIsNull;
+  }
+
   @Override
   public int getPositionCount() {
     return positionCount;
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/Column.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/Column.java
index 2e796a5f68..fe2af36fca 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/Column.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/Column.java
@@ -65,6 +65,41 @@ public interface Column {
     throw new UnsupportedOperationException(getClass().getName());
   }
 
+  /** Gets the boolean array. */
+  default boolean[] getBooleans() {
+    throw new UnsupportedOperationException(getClass().getName());
+  }
+
+  /** Gets the little endian int array. */
+  default int[] getInts() {
+    throw new UnsupportedOperationException(getClass().getName());
+  }
+
+  /** Gets the little endian long array. */
+  default long[] getLongs() {
+    throw new UnsupportedOperationException(getClass().getName());
+  }
+
+  /** Gets the float array. */
+  default float[] getFloats() {
+    throw new UnsupportedOperationException(getClass().getName());
+  }
+
+  /** Gets the double array. */
+  default double[] getDoubles() {
+    throw new UnsupportedOperationException(getClass().getName());
+  }
+
+  /** Gets the Binary list. */
+  default Binary[] getBinaries() {
+    throw new UnsupportedOperationException(getClass().getName());
+  }
+
+  /** Gets the Object array. */
+  default Object[] getObjects() {
+    throw new UnsupportedOperationException(getClass().getName());
+  }
+
   /** Gets a TsPrimitiveType at {@code position}. */
   default TsPrimitiveType getTsPrimitiveType(int position) {
     throw new UnsupportedOperationException(getClass().getName());
@@ -85,6 +120,9 @@ public interface Column {
    */
   boolean isNull(int position);
 
+  /** Returns the array to determine whether each position of the column is 
null or not. */
+  boolean[] isNull();
+
   /** Returns the number of positions in this block. */
   int getPositionCount();
 
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/DoubleColumn.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/DoubleColumn.java
index 25c97b188f..cac26a14cc 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/DoubleColumn.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/DoubleColumn.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
 import org.openjdk.jol.info.ClassLayout;
 
+import java.util.Arrays;
 import java.util.Optional;
 
 import static io.airlift.slice.SizeOf.sizeOf;
@@ -84,6 +85,11 @@ public class DoubleColumn implements Column {
     return values[position + arrayOffset];
   }
 
+  @Override
+  public double[] getDoubles() {
+    return values;
+  }
+
   @Override
   public Object getObject(int position) {
     return getDouble(position);
@@ -106,6 +112,16 @@ public class DoubleColumn implements Column {
     return valueIsNull != null && valueIsNull[position + arrayOffset];
   }
 
+  @Override
+  public boolean[] isNull() {
+    if (valueIsNull == null) {
+      boolean[] res = new boolean[positionCount];
+      Arrays.fill(res, false);
+      return res;
+    }
+    return valueIsNull;
+  }
+
   @Override
   public int getPositionCount() {
     return positionCount;
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/FloatColumn.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/FloatColumn.java
index b1831d15cb..12c1178691 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/FloatColumn.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/FloatColumn.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
 import org.openjdk.jol.info.ClassLayout;
 
+import java.util.Arrays;
 import java.util.Optional;
 
 import static io.airlift.slice.SizeOf.sizeOf;
@@ -83,6 +84,11 @@ public class FloatColumn implements Column {
     return values[position + arrayOffset];
   }
 
+  @Override
+  public float[] getFloats() {
+    return values;
+  }
+
   @Override
   public Object getObject(int position) {
     return getFloat(position);
@@ -105,6 +111,16 @@ public class FloatColumn implements Column {
     return valueIsNull != null && valueIsNull[position + arrayOffset];
   }
 
+  @Override
+  public boolean[] isNull() {
+    if (valueIsNull == null) {
+      boolean[] res = new boolean[positionCount];
+      Arrays.fill(res, false);
+      return res;
+    }
+    return valueIsNull;
+  }
+
   @Override
   public int getPositionCount() {
     return positionCount;
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/IntColumn.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/IntColumn.java
index f2556c75ed..5a6a71ab60 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/IntColumn.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/IntColumn.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
 import org.openjdk.jol.info.ClassLayout;
 
+import java.util.Arrays;
 import java.util.Optional;
 
 import static io.airlift.slice.SizeOf.sizeOf;
@@ -83,6 +84,11 @@ public class IntColumn implements Column {
     return values[position + arrayOffset];
   }
 
+  @Override
+  public int[] getInts() {
+    return values;
+  }
+
   @Override
   public Object getObject(int position) {
     return getInt(position);
@@ -105,6 +111,16 @@ public class IntColumn implements Column {
     return valueIsNull != null && valueIsNull[position + arrayOffset];
   }
 
+  @Override
+  public boolean[] isNull() {
+    if (valueIsNull == null) {
+      boolean[] res = new boolean[positionCount];
+      Arrays.fill(res, false);
+      return res;
+    }
+    return valueIsNull;
+  }
+
   @Override
   public int getPositionCount() {
     return positionCount;
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/LongColumn.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/LongColumn.java
index e4364a7072..3eb65f3d04 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/LongColumn.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/LongColumn.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
 import org.openjdk.jol.info.ClassLayout;
 
+import java.util.Arrays;
 import java.util.Optional;
 
 import static io.airlift.slice.SizeOf.sizeOf;
@@ -83,6 +84,11 @@ public class LongColumn implements Column {
     return values[position + arrayOffset];
   }
 
+  @Override
+  public long[] getLongs() {
+    return values;
+  }
+
   @Override
   public Object getObject(int position) {
     return getLong(position);
@@ -105,6 +111,16 @@ public class LongColumn implements Column {
     return valueIsNull != null && valueIsNull[position + arrayOffset];
   }
 
+  @Override
+  public boolean[] isNull() {
+    if (valueIsNull == null) {
+      boolean[] res = new boolean[positionCount];
+      Arrays.fill(res, false);
+      return res;
+    }
+    return valueIsNull;
+  }
+
   @Override
   public int getPositionCount() {
     return positionCount;
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/RunLengthEncodedColumn.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/RunLengthEncodedColumn.java
index 2dc69d43ec..0446782087 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/RunLengthEncodedColumn.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/RunLengthEncodedColumn.java
@@ -24,6 +24,8 @@ import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
 import org.openjdk.jol.info.ClassLayout;
 
+import java.util.Arrays;
+
 import static java.lang.String.format;
 import static java.util.Objects.requireNonNull;
 import static 
org.apache.iotdb.tsfile.read.common.block.column.ColumnUtil.checkValidRegion;
@@ -114,6 +116,55 @@ public class RunLengthEncodedColumn implements Column {
     return value.getObject(0);
   }
 
+  @Override
+  public boolean[] getBooleans() {
+    boolean[] res = new boolean[positionCount];
+    Arrays.fill(res, value.getBoolean(0));
+    return res;
+  }
+
+  @Override
+  public int[] getInts() {
+    int[] res = new int[positionCount];
+    Arrays.fill(res, value.getInt(0));
+    return res;
+  }
+
+  @Override
+  public long[] getLongs() {
+    long[] res = new long[positionCount];
+    Arrays.fill(res, value.getLong(0));
+    return res;
+  }
+
+  @Override
+  public float[] getFloats() {
+    float[] res = new float[positionCount];
+    Arrays.fill(res, value.getFloat(0));
+    return res;
+  }
+
+  @Override
+  public double[] getDoubles() {
+    double[] res = new double[positionCount];
+    Arrays.fill(res, value.getDouble(0));
+    return res;
+  }
+
+  @Override
+  public Binary[] getBinaries() {
+    Binary[] res = new Binary[positionCount];
+    Arrays.fill(res, value.getBinary(0));
+    return res;
+  }
+
+  @Override
+  public Object[] getObjects() {
+    Object[] res = new Object[positionCount];
+    Arrays.fill(res, value.getObject(0));
+    return res;
+  }
+
   @Override
   public TsPrimitiveType getTsPrimitiveType(int position) {
     checkReadablePosition(position);
@@ -131,6 +182,13 @@ public class RunLengthEncodedColumn implements Column {
     return value.isNull(0);
   }
 
+  @Override
+  public boolean[] isNull() {
+    boolean[] res = new boolean[positionCount];
+    Arrays.fill(res, value.isNull(0));
+    return res;
+  }
+
   @Override
   public int getPositionCount() {
     return positionCount;
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/TimeColumn.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/TimeColumn.java
index 5913a74547..af7ba86808 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/TimeColumn.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/TimeColumn.java
@@ -32,6 +32,7 @@ public class TimeColumn implements Column {
 
   private final int arrayOffset;
   private final int positionCount;
+
   private final long[] values;
 
   private final long retainedSizeInBytes;
@@ -92,6 +93,12 @@ public class TimeColumn implements Column {
     return false;
   }
 
+  @Override
+  public boolean[] isNull() {
+    // todo
+    return null;
+  }
+
   @Override
   public int getPositionCount() {
     return positionCount;
@@ -133,6 +140,10 @@ public class TimeColumn implements Column {
     return values[getPositionCount() + arrayOffset - 1];
   }
 
+  public long[] getTimes() {
+    return values;
+  }
+
   private void checkReadablePosition(int position) {
     if (position < 0 || position >= getPositionCount()) {
       throw new IllegalArgumentException("position is not valid");
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/AlignedChunkWriterImpl.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/AlignedChunkWriterImpl.java
index 102da36cb9..46d82d3a9f 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/AlignedChunkWriterImpl.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/AlignedChunkWriterImpl.java
@@ -23,6 +23,8 @@ import org.apache.iotdb.tsfile.encoding.encoder.Encoder;
 import org.apache.iotdb.tsfile.encoding.encoder.TSEncodingBuilder;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.iotdb.tsfile.read.common.block.column.Column;
+import org.apache.iotdb.tsfile.read.common.block.column.TimeColumn;
 import org.apache.iotdb.tsfile.utils.Binary;
 import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
@@ -40,6 +42,9 @@ public class AlignedChunkWriterImpl implements IChunkWriter {
   private final List<ValueChunkWriter> valueChunkWriterList;
   private int valueIndex;
 
+  // Used for batch writing
+  private long remainingPointsNumber;
+
   /** @param schema schema of this measurement */
   public AlignedChunkWriterImpl(VectorMeasurementSchema schema) {
     timeChunkWriter =
@@ -66,6 +71,7 @@ public class AlignedChunkWriterImpl implements IChunkWriter {
     }
 
     this.valueIndex = 0;
+    this.remainingPointsNumber = 
timeChunkWriter.getRemainingPointNumberForCurrentPage();
   }
 
   public AlignedChunkWriterImpl(List<IMeasurementSchema> schemaList) {
@@ -91,6 +97,8 @@ public class AlignedChunkWriterImpl implements IChunkWriter {
     }
 
     this.valueIndex = 0;
+
+    this.remainingPointsNumber = 
timeChunkWriter.getRemainingPointNumberForCurrentPage();
   }
 
   public void write(long time, int value, boolean isNull) {
@@ -156,6 +164,61 @@ public class AlignedChunkWriterImpl implements 
IChunkWriter {
     }
   }
 
+  public void write(TimeColumn timeColumn, Column[] valueColumns, int 
batchSize) {
+    if (remainingPointsNumber < batchSize) {
+      int pointsHasWritten = (int) remainingPointsNumber;
+      batchWrite(timeColumn, valueColumns, pointsHasWritten, 0);
+      batchWrite(timeColumn, valueColumns, batchSize - pointsHasWritten, 
pointsHasWritten);
+    } else {
+      batchWrite(timeColumn, valueColumns, batchSize, 0);
+    }
+  }
+
+  private void batchWrite(
+      TimeColumn timeColumn, Column[] valueColumns, int batchSize, int 
arrayOffset) {
+    valueIndex = 0;
+    long[] times = timeColumn.getTimes();
+
+    for (Column column : valueColumns) {
+      ValueChunkWriter chunkWriter = valueChunkWriterList.get(valueIndex++);
+      TSDataType tsDataType = chunkWriter.getDataType();
+      switch (tsDataType) {
+        case TEXT:
+          chunkWriter.write(times, column.getBinaries(), column.isNull(), 
batchSize, arrayOffset);
+          break;
+        case DOUBLE:
+          chunkWriter.write(times, column.getDoubles(), column.isNull(), 
batchSize, arrayOffset);
+          break;
+        case BOOLEAN:
+          chunkWriter.write(times, column.getBooleans(), column.isNull(), 
batchSize, arrayOffset);
+          break;
+        case INT64:
+          chunkWriter.write(times, column.getLongs(), column.isNull(), 
batchSize, arrayOffset);
+          break;
+        case INT32:
+          chunkWriter.write(times, column.getInts(), column.isNull(), 
batchSize, arrayOffset);
+          break;
+        case FLOAT:
+          chunkWriter.write(times, column.getFloats(), column.isNull(), 
batchSize, arrayOffset);
+          break;
+        default:
+          throw new UnsupportedOperationException("Unknown data type " + 
tsDataType);
+      }
+    }
+
+    write(times, batchSize, arrayOffset);
+  }
+
+  public void write(long[] time, int batchSize, int arrayOffset) {
+    valueIndex = 0;
+    timeChunkWriter.write(time, batchSize, arrayOffset);
+    if (checkPageSizeAndMayOpenANewPage()) {
+      writePageToPageBuffer();
+    }
+
+    remainingPointsNumber = 
timeChunkWriter.getRemainingPointNumberForCurrentPage();
+  }
+
   /**
    * check occupied memory size, if it exceeds the PageSize threshold, 
construct a page and put it
    * to pageBuffer
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/TimeChunkWriter.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/TimeChunkWriter.java
index b96f1a09ec..59a1fd465b 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/TimeChunkWriter.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/TimeChunkWriter.java
@@ -101,8 +101,8 @@ public class TimeChunkWriter {
     pageWriter.write(time);
   }
 
-  public void write(long[] timestamps, int batchSize) {
-    pageWriter.write(timestamps, batchSize);
+  public void write(long[] timestamps, int batchSize, int arrayOffset) {
+    pageWriter.write(timestamps, batchSize, arrayOffset);
   }
 
   /**
@@ -110,7 +110,7 @@ public class TimeChunkWriter {
    * to pageBuffer
    */
   public boolean checkPageSizeAndMayOpenANewPage() {
-    if (pageWriter.getPointNumber() == maxNumberOfPointsInPage) {
+    if (pageWriter.getPointNumber() >= maxNumberOfPointsInPage) {
       logger.debug("current line count reaches the upper bound, write page 
{}", measurementId);
       return true;
     } else if (pageWriter.getPointNumber()
@@ -136,6 +136,10 @@ public class TimeChunkWriter {
     return false;
   }
 
+  public long getRemainingPointNumberForCurrentPage() {
+    return maxNumberOfPointsInPage - pageWriter.getPointNumber();
+  }
+
   public void writePageToPageBuffer() {
     try {
       if (numOfPages == 0) { // record the firstPageStatistics
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ValueChunkWriter.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ValueChunkWriter.java
index e89edc4407..8b75269388 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ValueChunkWriter.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ValueChunkWriter.java
@@ -125,28 +125,28 @@ public class ValueChunkWriter {
     pageWriter.write(time, value, isNull);
   }
 
-  public void write(long[] timestamps, int[] values, int batchSize) {
-    pageWriter.write(timestamps, values, batchSize);
+  public void write(long[] timestamps, int[] values, boolean[] isNull, int 
batchSize, int pos) {
+    pageWriter.write(timestamps, values, isNull, batchSize, pos);
   }
 
-  public void write(long[] timestamps, long[] values, int batchSize) {
-    pageWriter.write(timestamps, values, batchSize);
+  public void write(long[] timestamps, long[] values, boolean[] isNull, int 
batchSize, int pos) {
+    pageWriter.write(timestamps, values, isNull, batchSize, pos);
   }
 
-  public void write(long[] timestamps, boolean[] values, int batchSize) {
-    pageWriter.write(timestamps, values, batchSize);
+  public void write(long[] timestamps, boolean[] values, boolean[] isNull, int 
batchSize, int pos) {
+    pageWriter.write(timestamps, values, isNull, batchSize, pos);
   }
 
-  public void write(long[] timestamps, float[] values, int batchSize) {
-    pageWriter.write(timestamps, values, batchSize);
+  public void write(long[] timestamps, float[] values, boolean[] isNull, int 
batchSize, int pos) {
+    pageWriter.write(timestamps, values, isNull, batchSize, pos);
   }
 
-  public void write(long[] timestamps, double[] values, int batchSize) {
-    pageWriter.write(timestamps, values, batchSize);
+  public void write(long[] timestamps, double[] values, boolean[] isNull, int 
batchSize, int pos) {
+    pageWriter.write(timestamps, values, isNull, batchSize, pos);
   }
 
-  public void write(long[] timestamps, Binary[] values, int batchSize) {
-    pageWriter.write(timestamps, values, batchSize);
+  public void write(long[] timestamps, Binary[] values, boolean[] isNull, int 
batchSize, int pos) {
+    pageWriter.write(timestamps, values, isNull, batchSize, pos);
   }
 
   public void writeEmptyPageToPageBuffer() {
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/page/TimePageWriter.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/page/TimePageWriter.java
index 1c668fc881..a865feeb62 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/page/TimePageWriter.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/page/TimePageWriter.java
@@ -67,11 +67,11 @@ public class TimePageWriter {
   }
 
   /** write time series into encoder */
-  public void write(long[] timestamps, int batchSize) {
-    for (int i = 0; i < batchSize; i++) {
+  public void write(long[] timestamps, int batchSize, int arrayOffset) {
+    for (int i = arrayOffset; i < batchSize + arrayOffset; i++) {
       timeEncoder.encode(timestamps[i], timeOut);
     }
-    statistics.update(timestamps, batchSize);
+    statistics.update(timestamps, batchSize, arrayOffset);
   }
 
   /** flush all data remained in encoders. */
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/page/ValuePageWriter.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/page/ValuePageWriter.java
index 988575aa2c..329a5c74d6 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/page/ValuePageWriter.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/page/ValuePageWriter.java
@@ -148,51 +148,75 @@ public class ValuePageWriter {
   }
 
   /** write time series into encoder */
-  public void write(long[] timestamps, boolean[] values, int batchSize) {
-    for (int i = 0; i < batchSize; i++) {
-      valueEncoder.encode(values[i], valueOut);
+  public void write(
+      long[] timestamps, boolean[] values, boolean[] isNull, int batchSize, 
int arrayOffset) {
+    for (int i = arrayOffset; i < batchSize + arrayOffset; i++) {
+      setBit(isNull[i]);
+      if (!isNull[i]) {
+        valueEncoder.encode(values[i], valueOut);
+        statistics.update(timestamps[i], values[i]);
+      }
     }
-    statistics.update(timestamps, values, batchSize);
   }
 
   /** write time series into encoder */
-  public void write(long[] timestamps, int[] values, int batchSize) {
-    for (int i = 0; i < batchSize; i++) {
-      valueEncoder.encode(values[i], valueOut);
+  public void write(
+      long[] timestamps, int[] values, boolean[] isNull, int batchSize, int 
arrayOffset) {
+    for (int i = arrayOffset; i < batchSize + arrayOffset; i++) {
+      setBit(isNull[i]);
+      if (!isNull[i]) {
+        valueEncoder.encode(values[i], valueOut);
+        statistics.update(timestamps[i], values[i]);
+      }
     }
-    statistics.update(timestamps, values, batchSize);
   }
 
   /** write time series into encoder */
-  public void write(long[] timestamps, long[] values, int batchSize) {
-    for (int i = 0; i < batchSize; i++) {
-      valueEncoder.encode(values[i], valueOut);
+  public void write(
+      long[] timestamps, long[] values, boolean[] isNull, int batchSize, int 
arrayOffset) {
+    for (int i = arrayOffset; i < batchSize + arrayOffset; i++) {
+      setBit(isNull[i]);
+      if (!isNull[i]) {
+        valueEncoder.encode(values[i], valueOut);
+        statistics.update(timestamps[i], values[i]);
+      }
     }
-    statistics.update(timestamps, values, batchSize);
   }
 
   /** write time series into encoder */
-  public void write(long[] timestamps, float[] values, int batchSize) {
-    for (int i = 0; i < batchSize; i++) {
-      valueEncoder.encode(values[i], valueOut);
+  public void write(
+      long[] timestamps, float[] values, boolean[] isNull, int batchSize, int 
arrayOffset) {
+    for (int i = arrayOffset; i < batchSize + arrayOffset; i++) {
+      setBit(isNull[i]);
+      if (!isNull[i]) {
+        valueEncoder.encode(values[i], valueOut);
+        statistics.update(timestamps[i], values[i]);
+      }
     }
-    statistics.update(timestamps, values, batchSize);
   }
 
   /** write time series into encoder */
-  public void write(long[] timestamps, double[] values, int batchSize) {
-    for (int i = 0; i < batchSize; i++) {
-      valueEncoder.encode(values[i], valueOut);
+  public void write(
+      long[] timestamps, double[] values, boolean[] isNull, int batchSize, int 
arrayOffset) {
+    for (int i = arrayOffset; i < batchSize + arrayOffset; i++) {
+      setBit(isNull[i]);
+      if (!isNull[i]) {
+        valueEncoder.encode(values[i], valueOut);
+        statistics.update(timestamps[i], values[i]);
+      }
     }
-    statistics.update(timestamps, values, batchSize);
   }
 
   /** write time series into encoder */
-  public void write(long[] timestamps, Binary[] values, int batchSize) {
-    for (int i = 0; i < batchSize; i++) {
-      valueEncoder.encode(values[i], valueOut);
+  public void write(
+      long[] timestamps, Binary[] values, boolean[] isNull, int batchSize, int 
arrayOffset) {
+    for (int i = arrayOffset; i < batchSize + arrayOffset; i++) {
+      setBit(isNull[i]);
+      if (!isNull[i]) {
+        valueEncoder.encode(values[i], valueOut);
+        statistics.update(timestamps[i], values[i]);
+      }
     }
-    statistics.update(timestamps, values, batchSize);
   }
 
   /** flush all data remained in encoders. */

Reply via email to