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. */