From dc3662bc4638bd8649694a97eb0492460705b0c5 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Wed, 16 Sep 2026 14:46:32 +0800 Subject: [PATCH 1/4] optimze --- .../apache/tsfile/tools/TabletBuilder.java | 44 ++- .../apache/tsfile/tools/ValueConverter.java | 17 ++ .../tsfile/tools/TabletBuilderTest.java | 32 +++ .../templates/PageDataReaderTemplate.ftl | 205 +++++++++++++ .../tsfile/read/reader/page/PageReader.java | 74 ++--- .../read/reader/page/ValuePageReader.java | 99 ++----- .../org/apache/tsfile/utils/TypeServices.java | 91 ++++++ .../chunk/AlignedChunkGroupWriterImpl.java | 12 +- .../tsfile/write/chunk/ValueChunkWriter.java | 5 + .../read/reader/PageDataBatchReaderTest.java | 272 ++++++++++++++++++ .../read/reader/PageDataTestSupport.java | 139 +++++++++ .../write/chunk/AlignedNullWriteTest.java | 98 +++++++ 12 files changed, 949 insertions(+), 139 deletions(-) create mode 100644 java/tsfile/src/main/codegen/templates/PageDataReaderTemplate.ftl create mode 100644 java/tsfile/src/test/java/org/apache/tsfile/read/reader/PageDataBatchReaderTest.java create mode 100644 java/tsfile/src/test/java/org/apache/tsfile/read/reader/PageDataTestSupport.java create mode 100644 java/tsfile/src/test/java/org/apache/tsfile/write/chunk/AlignedNullWriteTest.java diff --git a/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java b/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java index 001d0a1c2..e55019627 100644 --- a/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java +++ b/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java @@ -30,6 +30,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.function.Function; public class TabletBuilder { @@ -70,30 +71,43 @@ public Tablet build(SourceBatch batch) { Object timeValue = batch.getValue(row, timeColumnSourceIndex); long timestamp = timeConverter.convert(timeValue, importSchema.getTimePrecision()); tablet.addTimestamp(i, timestamp); + } - for (int col = 0; col < tableSchema.getColumnSchemas().size(); col++) { - IMeasurementSchema colSchema = tableSchema.getColumnSchemas().get(col); - String colName = colSchema.getMeasurementName(); + // SourceBatch and Tablet are columnar. Keep schema lookup and type dispatch outside the row + // loop. + for (int col = 0; col < tableSchema.getColumnSchemas().size(); col++) { + IMeasurementSchema colSchema = tableSchema.getColumnSchemas().get(col); + String colName = colSchema.getMeasurementName(); - if (tagDefaults.containsKey(colName)) { - tablet.addValue(colName, i, tagDefaults.get(colName)); - continue; + if (tagDefaults.containsKey(colName)) { + Object defaultValue = tagDefaults.get(colName); + for (int i = 0; i < rowCount; i++) { + tablet.addValue(colName, i, defaultValue); } + continue; + } - Integer srcIdx = sourceColumnIndex.get(colName); - if (srcIdx == null) { - continue; - } + Integer srcIdx = sourceColumnIndex.get(colName); + if (srcIdx == null) { + continue; + } - Object rawValue = batch.getValue(row, srcIdx); + Object[] sourceValues = batch.getColumn(srcIdx); + Function converter = null; + for (int i = 0; i < rowCount; i++) { + Object rawValue = sourceValues[sortedIndices[i]]; if (isNull(rawValue)) { continue; } - boolean isMeasurement = tableSchema.getColumnTypes().get(col) == ColumnCategory.FIELD; - Object converted = - ValueConverter.convert( - rawValue, colSchema.getType(), isMeasurement, importSchema.getTimePrecision()); + // Preserve the no-conversion behavior of empty/all-null columns. + if (converter == null) { + boolean isMeasurement = tableSchema.getColumnTypes().get(col) == ColumnCategory.FIELD; + converter = + ValueConverter.converterFor( + colSchema.getType(), isMeasurement, importSchema.getTimePrecision()); + } + Object converted = converter.apply(rawValue); tablet.addValue(colName, i, converted); } } diff --git a/java/tools/src/main/java/org/apache/tsfile/tools/ValueConverter.java b/java/tools/src/main/java/org/apache/tsfile/tools/ValueConverter.java index 75145f5e0..72861dc0f 100644 --- a/java/tools/src/main/java/org/apache/tsfile/tools/ValueConverter.java +++ b/java/tools/src/main/java/org/apache/tsfile/tools/ValueConverter.java @@ -29,6 +29,7 @@ import java.time.LocalDateTime; import java.time.ZoneId; import java.time.format.DateTimeParseException; +import java.util.function.Function; public class ValueConverter { @@ -150,6 +151,22 @@ private static Object fromString( .convert(value, isMeasurement, timePrecision); } + /** Resolve both input representations once for a column whose target type is fixed. */ + static Function converterFor( + TSDataType targetType, boolean isMeasurement, String timePrecision) { + Type type = Type.fromTsDataType(targetType); + StringConverter strings = FROM_STRING_SERVICE.call(type); + ObjectConverter objects = FROM_OBJECT_SERVICE.call(type); + return value -> { + if (value == null) { + return null; + } + return value instanceof String string + ? strings.convert(string, isMeasurement, timePrecision) + : objects.convert(value, isMeasurement, timePrecision); + }; + } + private static Object fromObject( Object value, TSDataType targetType, boolean isMeasurement, String timePrecision) { return FROM_OBJECT_SERVICE diff --git a/java/tools/src/test/java/org/apache/tsfile/tools/TabletBuilderTest.java b/java/tools/src/test/java/org/apache/tsfile/tools/TabletBuilderTest.java index 384841a25..5dcfbb3ef 100644 --- a/java/tools/src/test/java/org/apache/tsfile/tools/TabletBuilderTest.java +++ b/java/tools/src/test/java/org/apache/tsfile/tools/TabletBuilderTest.java @@ -35,6 +35,38 @@ public class TabletBuilderTest { + @Test + public void testColumnConversionKeepsSortedRowsAndNulls() { + ImportSchema schema = + buildSchema( + "test", + "time", + Collections.singletonList(new ImportSchema.TagColumn("region", "beijing")), + new ImportSchema.SourceColumn("time", TSDataType.INT64), + new ImportSchema.SourceColumn("number", TSDataType.INT32), + new ImportSchema.SourceColumn("flag", TSDataType.BOOLEAN)); + schema.setNullFormat("NULL"); + SourceBatch batch = + SourceBatch.fromRows( + Arrays.asList("time", "number", "flag"), + Arrays.asList( + new Object[] {30L, "30", "true"}, + new Object[] {10L, 10, false}, + new Object[] {20L, "NULL", ""})); + Tablet result = new TabletBuilder(schema, new TimeConverter("ms")).build(batch); + assertEquals(3, result.getRowSize()); + assertEquals(10L, result.getTimestamps()[0]); + assertEquals(10, result.getValue(0, 1)); + assertEquals(false, result.getValue(0, 2)); + assertTrue(result.isNull(1, 1)); + assertTrue(result.isNull(1, 2)); + assertEquals(30, result.getValue(2, 1)); + assertEquals(true, result.getValue(2, 2)); + for (int i = 0; i < 3; i++) { + assertTrue(result.getDeviceID(i).toString().contains("beijing")); + } + } + private ImportSchema buildSchema( String tableName, String timeCol, diff --git a/java/tsfile/src/main/codegen/templates/PageDataReaderTemplate.ftl b/java/tsfile/src/main/codegen/templates/PageDataReaderTemplate.ftl new file mode 100644 index 000000000..09db37370 --- /dev/null +++ b/java/tsfile/src/main/codegen/templates/PageDataReaderTemplate.ftl @@ -0,0 +1,205 @@ +<@pp.dropOutputFile /> +<#list [ + ["Boolean", "boolean", "Boolean", "Boolean", "TsBoolean"], + ["Int", "int", "Int", "Integer", "TsInt"], + ["Date", "int", "Int", "Integer", "TsInt"], + ["Long", "long", "Long", "Long", "TsLong"], + ["Float", "float", "Float", "Float", "TsFloat"], + ["Double", "double", "Double", "Double", "TsDouble"], + ["Binary", "Binary", "Binary", "Binary", "TsBinary"] +] as t> +<#assign name=t[0] primitive=t[1] suffix=t[2] filterSuffix=t[3] wrapper=t[4]> +<@pp.changeOutputFile name="/org/apache/tsfile/read/common/type/service/${name}PageDataReader.java" /> +/* + * 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.tsfile.read.common.type.service; + +import org.apache.tsfile.block.column.ColumnBuilder; +import org.apache.tsfile.encoding.decoder.Decoder; +<#if name == "Date"> +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.read.common.block.column.BinaryColumnBuilder; + +import org.apache.tsfile.read.common.BatchData; +import org.apache.tsfile.read.common.block.TsBlockBuilder; +import org.apache.tsfile.read.filter.basic.Filter; +import org.apache.tsfile.read.reader.series.PaginationController; +<#if name == "Binary"> +import org.apache.tsfile.utils.Binary; + +import org.apache.tsfile.utils.TsPrimitiveType; +import org.apache.tsfile.utils.TypeServices.PageDataBatchReader; + +import java.io.IOException; +import java.nio.ByteBuffer; +import java.util.function.LongPredicate; + +/** Generated from PageDataReaderTemplate.ftl. Dispatch once per batch, not once per value. */ +public final class ${name}PageDataReader implements PageDataBatchReader { + public static final ${name}PageDataReader INSTANCE = new ${name}PageDataReader(); + + private ${name}PageDataReader() {} + + @Override + public void readBatch(Decoder timeDecoder, ByteBuffer timeBuffer, Decoder decoder, + ByteBuffer buffer, Filter filter, BatchData output, boolean allSatisfy, + LongPredicate isDeleted) throws IOException { + while (timeDecoder.hasNext(timeBuffer)) { + long timestamp = timeDecoder.readLong(timeBuffer); + ${primitive} value = decoder.read${suffix}(buffer); + if (!isDeleted.test(timestamp) + && (allSatisfy || filter.satisfy${filterSuffix}(timestamp, value))) { + output.put${suffix}(timestamp, value); + } + } + } + + @Override + public void readAlignedBatch(long[] timestamps, byte[] bitmap, Decoder decoder, + ByteBuffer buffer, Filter filter, BatchData output, LongPredicate isDeleted) { + boolean allSatisfy = filter == null; + for (int i = 0; i < timestamps.length; i++) { + if ((bitmap[i / 8] & (0x80 >>> (i % 8))) == 0) { + continue; + } + long timestamp = timestamps[i]; + ${primitive} value = decoder.read${suffix}(buffer); + if (!isDeleted.test(timestamp) + && (allSatisfy || filter.satisfy${filterSuffix}(timestamp, value))) { + output.put${suffix}(timestamp, value); + } + } + } + + @Override + public long readBlock(Decoder timeDecoder, ByteBuffer timeBuffer, Decoder decoder, + ByteBuffer buffer, Filter filter, TsBlockBuilder builder, boolean allSatisfy, + LongPredicate isDeleted, PaginationController pagination) throws IOException { + ColumnBuilder times = builder.getTimeColumnBuilder(); + ColumnBuilder values = builder.getColumnBuilder(0); + long filtered = 0; + while (timeDecoder.hasNext(timeBuffer)) { + long timestamp = timeDecoder.readLong(timeBuffer); + // Consume both streams before deletion/filter/pagination checks, including the stop row. + ${primitive} value = decoder.read${suffix}(buffer); + if (isDeleted.test(timestamp)) { + continue; + } + if (!allSatisfy && !filter.satisfy${filterSuffix}(timestamp, value)) { + filtered++; + continue; + } + if (pagination.hasCurOffset()) { + pagination.consumeOffset(); + continue; + } + if (!pagination.hasCurLimit()) { + break; + } + times.writeLong(timestamp); + values.write${suffix}(value); + builder.declarePosition(); + pagination.consumeLimit(); + } + return filtered; + } + + @Override + public void readValues(long[] timestamps, byte[] bitmap, Decoder decoder, ByteBuffer buffer, + TsPrimitiveType[] output, LongPredicate isDeleted) { + for (int i = 0; i < output.length; i++) { + if ((bitmap[i / 8] & (0x80 >>> (i % 8))) == 0) { + continue; + } + ${primitive} value = decoder.read${suffix}(buffer); + if (!isDeleted.test(timestamps[i])) { + output[i] = new TsPrimitiveType.${wrapper}(value<#if name == "Date">, TSDataType.DATE); + } + } + } + + @Override + public void readColumn(int end, byte[] bitmap, Decoder decoder, ByteBuffer buffer, + ColumnBuilder builder, boolean[] keep, boolean[] deleted) { + <#if name == "Date"> + if (builder instanceof BinaryColumnBuilder binaryBuilder) { + <@selectedColumn dateBinary=true /> + return; + } + + <@selectedColumn dateBinary=false /> + } + + @Override + public void readColumn(int start, int end, byte[] bitmap, Decoder decoder, ByteBuffer buffer, + ColumnBuilder builder) { + // The encoded stream contains only present values; skipped rows still consume their values. + for (int i = 0; i < start; i++) { + if ((bitmap[i / 8] & (0x80 >>> (i % 8))) != 0) { + decoder.read${suffix}(buffer); + } + } + <#if name == "Date"> + if (builder instanceof BinaryColumnBuilder binaryBuilder) { + <@rangeColumn dateBinary=true /> + return; + } + + <@rangeColumn dateBinary=false /> + } +} + + +<#macro selectedColumn dateBinary> + for (int i = 0; i < end; i++) { + if ((bitmap[i / 8] & (0x80 >>> (i % 8))) == 0) { + if (keep[i]) { + builder.appendNull(); + } + continue; + } + // Decode even rows discarded by selection or deletion to keep the stream aligned. + ${primitive} value = decoder.read${suffix}(buffer); + if (keep[i]) { + if (deleted != null && deleted[i]) { + builder.appendNull(); + } else { + <#if dateBinary> + binaryBuilder.writeDate(value); + <#else> + builder.write${suffix}(value); + + } + } + } + + +<#macro rangeColumn dateBinary> + for (int i = start; i < end; i++) { + if ((bitmap[i / 8] & (0x80 >>> (i % 8))) == 0) { + builder.appendNull(); + } else { + <#if dateBinary> + binaryBuilder.writeDate(decoder.readInt(buffer)); + <#else> + builder.write${suffix}(decoder.read${suffix}(buffer)); + + } + } + diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/PageReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/PageReader.java index cdd45c4b2..ada833ee1 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/PageReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/PageReader.java @@ -35,9 +35,7 @@ import org.apache.tsfile.read.reader.series.PaginationController; import org.apache.tsfile.utils.ReadWriteForEncodingUtils; import org.apache.tsfile.utils.TypeServices; -import org.apache.tsfile.utils.TypeServices.PageDataBlockValueReader; -import org.apache.tsfile.utils.TypeServices.PageDataReadStatus; -import org.apache.tsfile.utils.TypeServices.PageDataValueReader; +import org.apache.tsfile.utils.TypeServices.PageDataBatchReader; import java.io.IOException; import java.io.Serializable; @@ -61,8 +59,7 @@ public class PageReader implements IPageReader { private final Decoder valueDecoder; // Reuse type-specific readers and the deletion predicate across repeated page reads. - private PageDataValueReader batchDataValueReader; - private PageDataBlockValueReader blockValueReader; + private PageDataBatchReader batchReader; private final LongPredicate deletePredicate = this::isDeleted; /** decoder for time column */ @@ -161,21 +158,16 @@ public BatchData getAllSatisfiedPageData(boolean ascending) throws IOException { uncompressDataIfNecessary(); BatchData pageData = BatchDataFactory.createBatchData(dataType, ascending, false); boolean allSatisfy = recordFilter == null || recordFilter.allSatisfy(this); - if (batchDataValueReader == null) { - batchDataValueReader = - TypeServices.READ_PAGE_VALUE_TO_BATCHDATA_SERVICE.call(Type.fromTsDataType(dataType)); - } - while (timeDecoder.hasNext(timeBuffer)) { - long timestamp = timeDecoder.readLong(timeBuffer); - batchDataValueReader.read( - valueDecoder, - valueBuffer, - recordFilter, - pageData, - timestamp, - allSatisfy, - deletePredicate); - } + getBatchReader() + .readBatch( + timeDecoder, + timeBuffer, + valueDecoder, + valueBuffer, + recordFilter, + pageData, + allSatisfy, + deletePredicate); return pageData.flip(); } @@ -195,36 +187,32 @@ public TsBlock getAllSatisfiedData(LongConsumer filterRowsRecorder) throws IOExc } builder = new TsBlockBuilder(initialExpectedEntries, Collections.singletonList(dataType)); - long allFilteredRows = 0; boolean allSatisfy = recordFilter == null || recordFilter.allSatisfy(this); - if (blockValueReader == null) { - blockValueReader = - TypeServices.READ_PAGE_VALUE_TO_TSBLOCK_SERVICE.call(Type.fromTsDataType(dataType)); - } - while (timeDecoder.hasNext(timeBuffer)) { - long timestamp = timeDecoder.readLong(timeBuffer); - PageDataReadStatus status = - blockValueReader.read( - valueDecoder, - valueBuffer, - recordFilter, - builder, - timestamp, - allSatisfy, - deletePredicate, - paginationController); - if (status == PageDataReadStatus.FILTERED) { - allFilteredRows++; - } else if (status == PageDataReadStatus.STOP) { - break; - } - } + long allFilteredRows = + getBatchReader() + .readBlock( + timeDecoder, + timeBuffer, + valueDecoder, + valueBuffer, + recordFilter, + builder, + allSatisfy, + deletePredicate, + paginationController); if (filterRowsRecorder != null && allFilteredRows > 0) { filterRowsRecorder.accept(allFilteredRows); } return builder.build(); } + private PageDataBatchReader getBatchReader() { + if (batchReader == null) { + batchReader = TypeServices.READ_PAGE_BATCH_SERVICE.call(Type.fromTsDataType(dataType)); + } + return batchReader; + } + @Override public Statistics getStatistics() { return pageHeader.getStatistics(); diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/ValuePageReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/ValuePageReader.java index 5151d25de..f3af223c1 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/ValuePageReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/ValuePageReader.java @@ -32,9 +32,8 @@ import org.apache.tsfile.utils.ReadWriteIOUtils; import org.apache.tsfile.utils.TsPrimitiveType; import org.apache.tsfile.utils.TypeServices; -import org.apache.tsfile.utils.TypeServices.PageDataColumnBuilderValueReader; +import org.apache.tsfile.utils.TypeServices.PageDataBatchReader; import org.apache.tsfile.utils.TypeServices.PageDataTsPrimitiveValueReader; -import org.apache.tsfile.utils.TypeServices.PageDataValueReader; import java.io.IOException; import java.io.Serializable; @@ -58,8 +57,7 @@ public class ValuePageReader { private PageDataTsPrimitiveValueReader valueReader; // Reuse readers for batch and column-builder APIs across repeated page reads. - private PageDataValueReader batchDataValueReader; - private PageDataColumnBuilderValueReader columnBuilderValueReader; + private PageDataBatchReader batchReader; // Reuse the bound predicate across page reads instead of creating a method reference per call. private final LongPredicate deletePredicate = this::isDeleted; @@ -128,19 +126,9 @@ public BatchData nextBatch(long[] timeBatch, boolean ascending, Filter filter) throws IOException { uncompressDataIfNecessary(); BatchData pageData = BatchDataFactory.createBatchData(dataType, ascending, false); - if (batchDataValueReader == null) { - batchDataValueReader = - TypeServices.READ_PAGE_VALUE_TO_BATCHDATA_SERVICE.call(Type.fromTsDataType(dataType)); - } - boolean allSatisfy = filter == null; - for (int i = 0; i < timeBatch.length; i++) { - if (((bitmap[i / 8] & 0xFF) & (MASK >>> (i % 8))) == 0) { - continue; - } - long timestamp = timeBatch[i]; - batchDataValueReader.read( - valueDecoder, valueBuffer, filter, pageData, timestamp, allSatisfy, deletePredicate); - } + getBatchReader() + .readAlignedBatch( + timeBatch, bitmap, valueDecoder, valueBuffer, filter, pageData, deletePredicate); return pageData.flip(); } @@ -167,17 +155,8 @@ public TsPrimitiveType[] nextValueBatch(long[] timeBatch) throws IOException { if (valueBuffer == null) { return valueBatch; } - if (valueReader == null) { - valueReader = - TypeServices.READ_PAGE_VALUE_TO_TSPRIMITIVETYPE_SERVICE.call( - Type.fromTsDataType(dataType)); - } - for (int i = 0; i < size; i++) { - if (((bitmap[i / 8] & 0xFF) & (MASK >>> (i % 8))) == 0) { - continue; - } - valueBatch[i] = valueReader.read(valueDecoder, valueBuffer, timeBatch[i], deletePredicate); - } + getBatchReader() + .readValues(timeBatch, bitmap, valueDecoder, valueBuffer, valueBatch, deletePredicate); return valueBatch; } @@ -193,20 +172,15 @@ public void writeColumnBuilderWithNextBatch( } return; } - if (columnBuilderValueReader == null) { - columnBuilderValueReader = - TypeServices.READ_PAGE_VALUE_TO_COLUMNBUILDER_SERVICE.call(Type.fromTsDataType(dataType)); - } - for (int i = 0; i < readEndIndex; i++) { - if (((bitmap[i / 8] & 0xFF) & (MASK >>> (i % 8))) == 0) { - if (keepCurrentRow[i]) { - columnBuilder.appendNull(); - } - continue; - } - columnBuilderValueReader.read( - valueDecoder, valueBuffer, columnBuilder, keepCurrentRow[i], isDeleted[i]); - } + getBatchReader() + .readColumn( + readEndIndex, + bitmap, + valueDecoder, + valueBuffer, + columnBuilder, + keepCurrentRow, + isDeleted); } public void writeColumnBuilderWithNextBatch( @@ -220,20 +194,9 @@ public void writeColumnBuilderWithNextBatch( } return; } - if (columnBuilderValueReader == null) { - columnBuilderValueReader = - TypeServices.READ_PAGE_VALUE_TO_COLUMNBUILDER_SERVICE.call(Type.fromTsDataType(dataType)); - } - for (int i = 0; i < readEndIndex; i++) { - if (((bitmap[i / 8] & 0xFF) & (MASK >>> (i % 8))) == 0) { - if (keepCurrentRow[i]) { - columnBuilder.appendNull(); - } - continue; - } - columnBuilderValueReader.read( - valueDecoder, valueBuffer, columnBuilder, keepCurrentRow[i], false); - } + getBatchReader() + .readColumn( + readEndIndex, bitmap, valueDecoder, valueBuffer, columnBuilder, keepCurrentRow, null); } public void writeColumnBuilderWithNextBatch( @@ -243,23 +206,15 @@ public void writeColumnBuilderWithNextBatch( columnBuilder.appendNull(readEndIndex - readStartIndex); return; } - if (columnBuilderValueReader == null) { - columnBuilderValueReader = - TypeServices.READ_PAGE_VALUE_TO_COLUMNBUILDER_SERVICE.call(Type.fromTsDataType(dataType)); - } - // skip useless data - for (int i = 0; i < readStartIndex; i++) { - if (((bitmap[i / 8] & 0xFF) & (MASK >>> (i % 8))) != 0) { - columnBuilderValueReader.read(valueDecoder, valueBuffer, columnBuilder, false, false); - } - } - for (int i = readStartIndex; i < readEndIndex; i++) { - if (((bitmap[i / 8] & 0xFF) & (MASK >>> (i % 8))) == 0) { - columnBuilder.appendNull(); - continue; - } - columnBuilderValueReader.read(valueDecoder, valueBuffer, columnBuilder, true, false); + getBatchReader() + .readColumn(readStartIndex, readEndIndex, bitmap, valueDecoder, valueBuffer, columnBuilder); + } + + private PageDataBatchReader getBatchReader() { + if (batchReader == null) { + batchReader = TypeServices.READ_PAGE_BATCH_SERVICE.call(Type.fromTsDataType(dataType)); } + return batchReader; } public Statistics getStatistics() { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/utils/TypeServices.java b/java/tsfile/src/main/java/org/apache/tsfile/utils/TypeServices.java index 34c5c16d6..01763e1af 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/utils/TypeServices.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/utils/TypeServices.java @@ -26,17 +26,108 @@ import org.apache.tsfile.read.common.BatchData; import org.apache.tsfile.read.common.block.TsBlockBuilder; import org.apache.tsfile.read.common.block.column.BinaryColumnBuilder; +import org.apache.tsfile.read.common.type.service.BinaryPageDataReader; +import org.apache.tsfile.read.common.type.service.BooleanPageDataReader; +import org.apache.tsfile.read.common.type.service.DatePageDataReader; +import org.apache.tsfile.read.common.type.service.DoublePageDataReader; +import org.apache.tsfile.read.common.type.service.FloatPageDataReader; +import org.apache.tsfile.read.common.type.service.IntPageDataReader; +import org.apache.tsfile.read.common.type.service.LongPageDataReader; import org.apache.tsfile.read.common.type.service.TypeService; import org.apache.tsfile.read.filter.basic.Filter; import org.apache.tsfile.read.reader.series.PaginationController; import org.apache.tsfile.write.UnSupportedDataTypeException; import org.apache.tsfile.write.chunk.ValueChunkWriter; +import java.io.IOException; import java.nio.ByteBuffer; import java.util.function.LongPredicate; public final class TypeServices { + /** Type-specialized loops avoid a polymorphic reader invocation for every decoded value. */ + public static final TypeService READ_PAGE_BATCH_SERVICE = + type -> + switch (type.getTypeEnum()) { + case BOOLEAN -> BooleanPageDataReader.INSTANCE; + case INT32 -> IntPageDataReader.INSTANCE; + case DATE -> DatePageDataReader.INSTANCE; + case INT64, TIMESTAMP -> LongPageDataReader.INSTANCE; + case FLOAT -> FloatPageDataReader.INSTANCE; + case DOUBLE -> DoublePageDataReader.INSTANCE; + case TEXT, BLOB, STRING, OBJECT -> BinaryPageDataReader.INSTANCE; + case ROW, UNKNOWN, VECTOR -> + throw new UnSupportedDataTypeException(String.valueOf(type.getTypeEnum())) + .setChecked(true); + }; + + static { + READ_PAGE_BATCH_SERVICE.check(); + } + + /** Batch entry points; implementations own the loop and directly use typed decoder methods. */ + public interface PageDataBatchReader { + void readBatch( + Decoder timeDecoder, + ByteBuffer timeBuffer, + Decoder decoder, + ByteBuffer buffer, + Filter filter, + BatchData output, + boolean allSatisfy, + LongPredicate isDeleted) + throws IOException; + + void readAlignedBatch( + long[] timestamps, + byte[] bitmap, + Decoder decoder, + ByteBuffer buffer, + Filter filter, + BatchData output, + LongPredicate isDeleted); + + /** Returns the number of rows rejected by the record filter, excluding deleted rows. */ + long readBlock( + Decoder timeDecoder, + ByteBuffer timeBuffer, + Decoder decoder, + ByteBuffer buffer, + Filter filter, + TsBlockBuilder builder, + boolean allSatisfy, + LongPredicate isDeleted, + PaginationController pagination) + throws IOException; + + void readValues( + long[] timestamps, + byte[] bitmap, + Decoder decoder, + ByteBuffer buffer, + TsPrimitiveType[] output, + LongPredicate isDeleted); + + /** A null deletion mask means that no rows are deleted. */ + void readColumn( + int end, + byte[] bitmap, + Decoder decoder, + ByteBuffer buffer, + ColumnBuilder builder, + boolean[] keep, + boolean[] deleted); + + /** Consumes present values before start, then appends the range [start, end). */ + void readColumn( + int start, + int end, + byte[] bitmap, + Decoder decoder, + ByteBuffer buffer, + ColumnBuilder builder); + } + // Page value decoding services for BatchData and TsBlock outputs. public static final TypeService READ_PAGE_VALUE_TO_BATCHDATA_SERVICE = diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java index 8896b58a4..6b66792ef 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java @@ -32,8 +32,6 @@ import org.apache.tsfile.file.metadata.enums.TSEncoding; import org.apache.tsfile.i18n.Messages; import org.apache.tsfile.read.common.type.Type; -import org.apache.tsfile.utils.TypeServices; -import org.apache.tsfile.utils.TypeServices.EmptyValueChunkWriter; import org.apache.tsfile.write.record.Tablet; import org.apache.tsfile.write.record.datapoint.DataPoint; import org.apache.tsfile.write.schema.IMeasurementSchema; @@ -285,17 +283,13 @@ public void tryToAddEmptyPageAndData(ValueChunkWriter valueChunkWriter) throws I } // add empty data of currentPage - for (long i = 0; i < timeChunkWriter.getPageWriter().getStatistics().getCount(); i++) { - valueChunkWriter.write(0, 0, true); - } + valueChunkWriter.writeNull(timeChunkWriter.getPageWriter().getStatistics().getCount()); } private void writeEmptyDataInOneRow(List valueChunkWriterList) { for (ValueChunkWriter valueChunkWriter : valueChunkWriterList) { - EmptyValueChunkWriter emptyValueWriter = - TypeServices.WRITE_EMPTY_VALUE_TO_CHUNK_SERVICE.call( - Type.fromTsDataType(valueChunkWriter.getDataType())); - emptyValueWriter.write(valueChunkWriter); + // Keep row-wise page-size checks and synchronized page boundaries across all columns. + valueChunkWriter.writeNull(1); } } diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/ValueChunkWriter.java b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/ValueChunkWriter.java index 9ade788c1..382df6e80 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/ValueChunkWriter.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/ValueChunkWriter.java @@ -164,6 +164,11 @@ public void write(long time, Binary value, boolean isNull) { pageWriter.write(time, value, isNull); } + /** Appends nulls to the current page without type dispatch, encoding, or statistics updates. */ + public void writeNull(int count) { + pageWriter.writeNull(count); + } + public void write(long[] timestamps, int[] values, boolean[] isNull, int batchSize, int pos) { pageWriter.write(timestamps, values, isNull, batchSize, pos); } diff --git a/java/tsfile/src/test/java/org/apache/tsfile/read/reader/PageDataBatchReaderTest.java b/java/tsfile/src/test/java/org/apache/tsfile/read/reader/PageDataBatchReaderTest.java new file mode 100644 index 000000000..4a0c31270 --- /dev/null +++ b/java/tsfile/src/test/java/org/apache/tsfile/read/reader/PageDataBatchReaderTest.java @@ -0,0 +1,272 @@ +/* + * 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.tsfile.read.reader; + +import org.apache.tsfile.block.column.Column; +import org.apache.tsfile.block.column.ColumnBuilder; +import org.apache.tsfile.encoding.decoder.PlainDecoder; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.read.common.BatchData; +import org.apache.tsfile.read.common.BatchDataFactory; +import org.apache.tsfile.read.common.block.TsBlock; +import org.apache.tsfile.read.common.block.TsBlockBuilder; +import org.apache.tsfile.read.common.block.column.BinaryColumnBuilder; +import org.apache.tsfile.read.common.type.Type; +import org.apache.tsfile.read.filter.basic.Filter; +import org.apache.tsfile.read.filter.factory.TimeFilterApi; +import org.apache.tsfile.read.reader.page.ValuePageReader; +import org.apache.tsfile.read.reader.series.PaginationController; +import org.apache.tsfile.utils.TsPrimitiveType; +import org.apache.tsfile.utils.TypeServices; +import org.apache.tsfile.utils.TypeServices.PageDataReadStatus; + +import org.junit.Test; + +import java.nio.ByteBuffer; +import java.util.Arrays; +import java.util.Collections; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; + +public class PageDataBatchReaderTest { + private static final TSDataType[] TYPES = { + TSDataType.BOOLEAN, + TSDataType.INT32, + TSDataType.DATE, + TSDataType.INT64, + TSDataType.TIMESTAMP, + TSDataType.FLOAT, + TSDataType.DOUBLE, + TSDataType.TEXT, + TSDataType.STRING, + TSDataType.BLOB, + TSDataType.OBJECT + }; + + @Test + public void columnsMatchScalarReaders() throws Exception { + for (TSDataType type : TYPES) { + for (int size : new int[] {0, 1, 7, 8, 9, 65}) { + PageDataTestSupport f = new PageDataTestSupport(type, size, true); + for (int mode = 0; mode < 4; mode++) { + for (boolean binaryDate : new boolean[] {false, true}) { + if (binaryDate && type != TSDataType.DATE) { + continue; + } + ColumnBuilder expected = builder(type, size, binaryDate); + ColumnBuilder actual = builder(type, size, binaryDate); + ByteBuffer input = ByteBuffer.wrap(f.values); + PlainDecoder decoder = new PlainDecoder(); + var scalar = + TypeServices.READ_PAGE_VALUE_TO_COLUMNBUILDER_SERVICE.call( + Type.fromTsDataType(type)); + boolean[] keep = f.keep.clone(); + if (mode == 3) { + Arrays.fill(keep, false); + } + int start = mode == 2 ? size / 3 : 0; + for (int i = 0; i < size; i++) { + boolean selected = mode == 2 ? i >= start : keep[i]; + if ((f.bitmap[i / 8] & (0x80 >>> (i % 8))) == 0) { + if (selected) { + expected.appendNull(); + } + } else { + scalar.read(decoder, input, expected, selected, mode == 0 && f.deleted[i]); + } + } + ValuePageReader reader = f.alignedReader(false); + if (mode == 0) { + reader.writeColumnBuilderWithNextBatch(size, actual, keep, f.deleted); + } else if (mode == 2) { + reader.writeColumnBuilderWithNextBatch(start, size, actual); + } else { + reader.writeColumnBuilderWithNextBatch(size, actual, keep); + } + assertColumn(expected.build(), actual.build()); + } + } + } + } + } + + @Test + public void batchesAndPrimitiveArraysMatchScalarReaders() throws Exception { + for (TSDataType type : TYPES) { + PageDataTestSupport f = new PageDataTestSupport(type, 65, true); + for (boolean aligned : new boolean[] {false, true}) { + for (boolean ascending : new boolean[] {false, true}) { + for (Filter filter : new Filter[] {null, TimeFilterApi.gt(10)}) { + BatchData expected = BatchDataFactory.createBatchData(type, ascending, false); + ByteBuffer input = ByteBuffer.wrap(aligned ? f.values : f.allValues); + PlainDecoder decoder = new PlainDecoder(); + var scalar = + TypeServices.READ_PAGE_VALUE_TO_BATCHDATA_SERVICE.call(Type.fromTsDataType(type)); + for (int i = 0; i < f.size; i++) { + if (!aligned || (f.bitmap[i / 8] & (0x80 >>> (i % 8))) != 0) { + scalar.read( + decoder, input, filter, expected, i, filter == null, t -> f.deleted[(int) t]); + } + } + expected.flip(); + BatchData actual = + aligned + ? f.alignedReader(true).nextBatch(f.times, ascending, filter) + : f.pageReader(filter, true).getAllSatisfiedPageData(ascending); + while (expected.hasCurrent()) { + assertTrue(actual.hasCurrent()); + assertEquals(expected.currentTime(), actual.currentTime()); + assertEquals(expected.currentValue(), actual.currentValue()); + expected.next(); + actual.next(); + } + assertFalse(actual.hasCurrent()); + } + } + } + TsPrimitiveType[] actual = f.alignedReader(true).nextValueBatch(f.times); + ValuePageReader scalarReader = f.alignedReader(true); + for (int i = 0; i < f.size; i++) { + TsPrimitiveType expected = scalarReader.nextValue(i, i); + if (expected == null) { + assertNull(actual[i]); + } else { + assertEquals(expected.getDataType(), actual[i].getDataType()); + assertEquals(expected.getValue(), actual[i].getValue()); + } + } + } + } + + @Test + public void blockPaginationPreservesFilterCountsAndDecoderPositions() throws Exception { + for (TSDataType type : TYPES) { + PageDataTestSupport f = new PageDataTestSupport(type, 65, false); + for (Filter filter : new Filter[] {null, TimeFilterApi.gt(10), TimeFilterApi.gt(100)}) { + for (long limit : new long[] {0, 1, 8, 100}) { + for (long offset : new long[] {0, 5, 100}) { + var scalar = + TypeServices.READ_PAGE_VALUE_TO_TSBLOCK_SERVICE.call(Type.fromTsDataType(type)); + var batch = TypeServices.READ_PAGE_BATCH_SERVICE.call(Type.fromTsDataType(type)); + ByteBuffer expectedTimes = ByteBuffer.wrap(f.timeBytes); + ByteBuffer expectedValues = ByteBuffer.wrap(f.allValues); + ByteBuffer actualTimes = ByteBuffer.wrap(f.timeBytes); + ByteBuffer actualValues = ByteBuffer.wrap(f.allValues); + TsBlockBuilder expected = new TsBlockBuilder(Collections.singletonList(type)); + TsBlockBuilder actual = new TsBlockBuilder(Collections.singletonList(type)); + PaginationController p1 = new PaginationController(limit, offset); + PaginationController p2 = new PaginationController(limit, offset); + PlainDecoder decoder = new PlainDecoder(); + long filtered = 0; + while (expectedTimes.hasRemaining()) { + long timestamp = decoder.readLong(expectedTimes); + PageDataReadStatus status = + scalar.read( + decoder, + expectedValues, + filter, + expected, + timestamp, + filter == null, + t -> t >= 20 && t <= 30, + p1); + if (status == PageDataReadStatus.FILTERED) { + filtered++; + } else if (status == PageDataReadStatus.STOP) { + break; + } + } + long actualFiltered = + batch.readBlock( + new PlainDecoder(), + actualTimes, + new PlainDecoder(), + actualValues, + filter, + actual, + filter == null, + t -> t >= 20 && t <= 30, + p2); + assertEquals(filtered, actualFiltered); + assertEquals(expectedTimes.position(), actualTimes.position()); + assertEquals(expectedValues.position(), actualValues.position()); + assertEquals(p1.getCurLimit(), p2.getCurLimit()); + assertEquals(p1.getCurOffset(), p2.getCurOffset()); + TsBlock expectedBlock = expected.build(); + TsBlock actualBlock = actual.build(); + assertColumn(expectedBlock.getTimeColumn(), actualBlock.getTimeColumn()); + assertColumn(expectedBlock.getColumn(0), actualBlock.getColumn(0)); + } + } + } + } + } + + private static ColumnBuilder builder(TSDataType type, int size, boolean binaryDate) { + return binaryDate + ? new BinaryColumnBuilder(null, size) + : Type.fromTsDataType(type).createColumnBuilder(size); + } + + @Test + public void discardedRowsConsumeValuesAndEmptyPagesAppendNulls() throws Exception { + for (TSDataType type : TYPES) { + PageDataTestSupport f = new PageDataTestSupport(type, 65, true); + ValuePageReader reader = f.alignedReader(false); + ValuePageReader scalar = f.alignedReader(false); + ColumnBuilder output = builder(type, 65, false); + reader.writeColumnBuilderWithNextBatch(32, output, new boolean[32]); + assertEquals(0, output.build().getPositionCount()); + for (int i = 0; i < 32; i++) { + scalar.nextValue(i, i); + } + for (int i = 32; i < 65; i++) { + TsPrimitiveType expected = scalar.nextValue(i, i); + TsPrimitiveType actual = reader.nextValue(i, i); + assertEquals( + expected == null ? null : expected.getValue(), + actual == null ? null : actual.getValue()); + } + ValuePageReader empty = + new ValuePageReader(f.header, (ByteBuffer) null, type, new PlainDecoder()); + ColumnBuilder nulls = builder(type, 5, false); + empty.writeColumnBuilderWithNextBatch(3, nulls, new boolean[] {true, false, true}); + empty.writeColumnBuilderWithNextBatch(1, 4, nulls); + Column result = nulls.build(); + assertEquals(5, result.getPositionCount()); + for (int i = 0; i < 5; i++) { + assertTrue(result.isNull(i)); + } + assertEquals(0, empty.nextValueBatch(new long[0]).length); + } + } + + private static void assertColumn(Column expected, Column actual) { + assertEquals(expected.getPositionCount(), actual.getPositionCount()); + for (int i = 0; i < expected.getPositionCount(); i++) { + assertEquals(expected.isNull(i), actual.isNull(i)); + if (!expected.isNull(i)) { + assertEquals(expected.getObject(i), actual.getObject(i)); + } + } + } +} diff --git a/java/tsfile/src/test/java/org/apache/tsfile/read/reader/PageDataTestSupport.java b/java/tsfile/src/test/java/org/apache/tsfile/read/reader/PageDataTestSupport.java new file mode 100644 index 000000000..2d4066da4 --- /dev/null +++ b/java/tsfile/src/test/java/org/apache/tsfile/read/reader/PageDataTestSupport.java @@ -0,0 +1,139 @@ +/* + * 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.tsfile.read.reader; + +import org.apache.tsfile.encoding.decoder.PlainDecoder; +import org.apache.tsfile.encoding.encoder.PlainEncoder; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.file.header.PageHeader; +import org.apache.tsfile.file.metadata.statistics.Statistics; +import org.apache.tsfile.read.common.TimeRange; +import org.apache.tsfile.read.filter.basic.Filter; +import org.apache.tsfile.read.reader.page.PageReader; +import org.apache.tsfile.read.reader.page.ValuePageReader; +import org.apache.tsfile.utils.Binary; +import org.apache.tsfile.utils.ReadWriteForEncodingUtils; + +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.util.Collections; + +/** Deterministic encoded pages shared by regression tests and the standalone JMH benchmark. */ +public class PageDataTestSupport { + public final TSDataType type; + public final int size; + public final byte[] bitmap; + public final long[] times; + public final boolean[] keep; + public final boolean[] deleted; + public final byte[] values; + public final byte[] allValues; + public final byte[] timeBytes; + public final byte[] aligned; + public final byte[] nonAligned; + public final PageHeader header; + + public PageDataTestSupport(TSDataType type, int size, boolean sparse) throws IOException { + this.type = type; + this.size = size; + bitmap = new byte[(size + 7) / 8]; + times = new long[size]; + keep = new boolean[size]; + deleted = new boolean[size]; + ByteArrayOutputStream valueOut = new ByteArrayOutputStream(); + ByteArrayOutputStream allValueOut = new ByteArrayOutputStream(); + ByteArrayOutputStream timeOut = new ByteArrayOutputStream(); + PlainEncoder encoder = new PlainEncoder(type, 128); + PlainEncoder timeEncoder = new PlainEncoder(TSDataType.INT64, 0); + for (int i = 0; i < size; i++) { + times[i] = i; + timeEncoder.encode((long) i, timeOut); + keep[i] = !sparse || i % 3 != 0; + deleted[i] = sparse && i >= size / 3 && i < size / 2; + encode(encoder, allValueOut, type, i); + if (!sparse || i % 5 != 0) { + bitmap[i / 8] |= (byte) (0x80 >>> (i % 8)); + encode(encoder, valueOut, type, i); + } + } + values = valueOut.toByteArray(); + allValues = allValueOut.toByteArray(); + timeBytes = timeOut.toByteArray(); + aligned = + ByteBuffer.allocate(4 + bitmap.length + values.length) + .putInt(size) + .put(bitmap) + .put(values) + .array(); + ByteArrayOutputStream page = new ByteArrayOutputStream(); + ReadWriteForEncodingUtils.writeUnsignedVarInt(timeBytes.length, page); + page.write(timeBytes); + page.write(allValues); + nonAligned = page.toByteArray(); + Statistics statistics = Statistics.getStatsByType(type); + statistics.setCount(size); + statistics.setStartTime(0); + statistics.setEndTime(size - 1L); + statistics.setEmpty(false); + header = new PageHeader(nonAligned.length, nonAligned.length, statistics); + } + + public ValuePageReader alignedReader(boolean withDeletion) { + ValuePageReader reader = + new ValuePageReader(header, ByteBuffer.wrap(aligned), type, new PlainDecoder()); + if (withDeletion) { + reader.setDeleteIntervalList( + Collections.singletonList(new TimeRange(size / 3, size / 2 - 1))); + } + return reader; + } + + public PageReader pageReader(Filter filter, boolean withDeletion) { + PageReader reader = + new PageReader( + header, + ByteBuffer.wrap(nonAligned), + type, + new PlainDecoder(), + new PlainDecoder(), + filter); + if (withDeletion) { + reader.setDeleteIntervalList( + Collections.singletonList(new TimeRange(size / 3, size / 2 - 1))); + } + return reader; + } + + private static void encode( + PlainEncoder encoder, ByteArrayOutputStream out, TSDataType type, int i) { + switch (type) { + case BOOLEAN -> encoder.encode(i % 2 == 0, out); + case INT32 -> encoder.encode(i * 13 - 50, out); + case DATE -> encoder.encode(20240101 + i % 28, out); + case INT64, TIMESTAMP -> encoder.encode(i * 1000003L, out); + case FLOAT -> encoder.encode(i * 0.25f, out); + case DOUBLE -> encoder.encode(i * 0.125, out); + case TEXT, STRING, BLOB, OBJECT -> + encoder.encode(new Binary("value-" + i, StandardCharsets.UTF_8), out); + case VECTOR, UNKNOWN -> throw new AssertionError(type); + } + } +} diff --git a/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/AlignedNullWriteTest.java b/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/AlignedNullWriteTest.java new file mode 100644 index 000000000..8e6eebe36 --- /dev/null +++ b/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/AlignedNullWriteTest.java @@ -0,0 +1,98 @@ +/* + * 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.tsfile.write.chunk; + +import org.apache.tsfile.common.conf.TSFileConfig; +import org.apache.tsfile.common.conf.TSFileDescriptor; +import org.apache.tsfile.encoding.decoder.PlainDecoder; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.file.metadata.IDeviceID; +import org.apache.tsfile.file.metadata.enums.CompressionType; +import org.apache.tsfile.file.metadata.enums.TSEncoding; +import org.apache.tsfile.write.record.datapoint.IntDataPoint; +import org.apache.tsfile.write.schema.MeasurementSchema; + +import org.junit.Test; + +import java.nio.ByteBuffer; +import java.util.Collections; + +import static org.junit.Assert.assertEquals; + +public class AlignedNullWriteTest { + @Test + public void lateColumnsAndMissingRowsKeepPageBoundaries() throws Exception { + TSFileConfig config = TSFileDescriptor.getInstance().getConfig(); + int oldLimit = config.getMaxNumberOfPointsInPage(); + try { + config.setMaxNumberOfPointsInPage(8); + AlignedChunkGroupWriterImpl group = + new AlignedChunkGroupWriterImpl(IDeviceID.Factory.DEFAULT_FACTORY.create("root.nulls")); + for (int i = 0; i < 11; i++) { + group.write(i, Collections.emptyList()); + } + TSDataType[] types = { + TSDataType.BOOLEAN, + TSDataType.INT32, + TSDataType.DATE, + TSDataType.INT64, + TSDataType.TIMESTAMP, + TSDataType.FLOAT, + TSDataType.DOUBLE, + TSDataType.TEXT, + TSDataType.STRING, + TSDataType.BLOB, + TSDataType.OBJECT + }; + // Each newly added column must catch up one sealed page and three rows of the current page. + for (TSDataType type : types) { + ValueChunkWriter writer = + group.tryToAddSeriesWriterInternal( + new MeasurementSchema( + type.name(), type, TSEncoding.PLAIN, CompressionType.UNCOMPRESSED)); + assertEquals(1, writer.getNumOfPages()); + assertEquals(3, writer.getPageWriter().getSize()); + assertEquals(0, writer.getPageWriter().getStatistics().getCount()); + } + for (int i = 11; i < 16; i++) { + group.write(i, Collections.emptyList()); + } + assertEquals(2, group.timeChunkWriter.getNumOfPages()); + for (ValueChunkWriter writer : group.valueChunkWriterMap.values()) { + assertEquals(2, writer.getNumOfPages()); + assertEquals(0, writer.getPageWriter().getSize()); + } + group.write(16, Collections.singletonList(new IntDataPoint("INT32", 42))); + for (ValueChunkWriter writer : group.valueChunkWriterMap.values()) { + assertEquals(1, writer.getPageWriter().getSize()); + boolean present = writer.getDataType() == TSDataType.INT32; + assertEquals(present ? 1 : 0, writer.getPageWriter().getStatistics().getCount()); + ByteBuffer data = writer.getPageWriter().getUncompressedBytes(); + assertEquals(1, data.getInt()); + assertEquals(present ? (byte) 0x80 : 0, data.get()); + if (present) { + assertEquals(42, new PlainDecoder().readInt(data)); + } + assertEquals(0, data.remaining()); + } + } finally { + config.setMaxNumberOfPointsInPage(oldLimit); + } + } +} From 9d2e028b0e63ee2f1892f031c7f8d99197581d59 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Wed, 16 Sep 2026 17:52:51 +0800 Subject: [PATCH 2/4] tmp --- .../tsfile/tools/ArrowSourceReader.java | 123 +++++++++++------- .../tsfile/tools/ParquetSourceReader.java | 91 ++++++------- .../apache/tsfile/tools/TabletBuilder.java | 11 +- .../chunk/AlignedChunkGroupWriterImpl.java | 15 ++- .../chunk/NonAlignedChunkGroupWriterImpl.java | 41 ++++-- .../write/chunk/AlignedNullWriteTest.java | 57 ++++++++ 6 files changed, 222 insertions(+), 116 deletions(-) diff --git a/java/tools/src/main/java/org/apache/tsfile/tools/ArrowSourceReader.java b/java/tools/src/main/java/org/apache/tsfile/tools/ArrowSourceReader.java index e23222d80..3f439224f 100644 --- a/java/tools/src/main/java/org/apache/tsfile/tools/ArrowSourceReader.java +++ b/java/tools/src/main/java/org/apache/tsfile/tools/ArrowSourceReader.java @@ -182,23 +182,15 @@ public SourceBatch readBatch() { } int numCols = schemaColumnNames.size(); - List rows = new ArrayList<>(rowCount); - - for (int r = 0; r < rowCount; r++) { - Object[] row = new Object[numCols]; - for (int c = 0; c < numCols; c++) { - String colName = schemaColumnNames.get(c); - FieldVector vec = vectorMap.get(colName); - if (vec == null || vec.isNull(r)) { - row[c] = null; - } else { - row[c] = extractValue(vec, r); - } + Object[][] columns = new Object[numCols][rowCount]; + // Arrow and SourceBatch are columnar: resolve each vector once and avoid a row transpose. + for (int c = 0; c < numCols; c++) { + FieldVector vec = vectorMap.get(schemaColumnNames.get(c)); + if (vec != null) { + readColumn(vec, columns[c]); } - rows.add(row); } - - return SourceBatch.fromRows(schemaColumnNames, rows); + return new SourceBatch(schemaColumnNames.toArray(new String[0]), columns, rowCount); } catch (IOException e) { LOGGER.error(Messages.format("log.tools.arrow_read_error", sourceFile.getAbsolutePath()), e); exhausted = true; @@ -293,40 +285,75 @@ private List getSchemaColumnNames() { return names; } - private Object extractValue(FieldVector vec, int row) { - // Date / Timestamp checks must come BEFORE the BigIntVector/IntVector branches: although - // they hold int/long underneath, DateDayVector / TimeStampVector do NOT extend - // IntVector / BigIntVector, so without these branches Date columns fall through to the - // generic getObject().toString() path and produce strings that don't match TSDataType.DATE. - if (vec instanceof DateDayVector) { - // Days since 1970-01-01. ValueConverter.toLocalDate handles Integer → LocalDate. - return ((DateDayVector) vec).get(row); - } else if (vec instanceof DateMilliVector) { - // Millis since 1970-01-01; collapse to date. - long millis = ((DateMilliVector) vec).get(row); - return LocalDate.ofEpochDay(Math.floorDiv(millis, 86_400_000L)); - } else if (vec instanceof TimeStampVector) { - // Long in the vector's native precision; matches the precision detected by - // detectTimestampPrecision() and stored on the schema. - return ((TimeStampVector) vec).get(row); - } else if (vec instanceof BigIntVector) { - return ((BigIntVector) vec).get(row); - } else if (vec instanceof IntVector) { - return ((IntVector) vec).get(row); - } else if (vec instanceof Float4Vector) { - return ((Float4Vector) vec).get(row); - } else if (vec instanceof Float8Vector) { - return ((Float8Vector) vec).get(row); - } else if (vec instanceof BitVector) { - return ((BitVector) vec).get(row) != 0; - } else if (vec instanceof VarCharVector) { - byte[] bytes = ((VarCharVector) vec).get(row); - return new String(bytes, StandardCharsets.UTF_8); - } else if (vec instanceof VarBinaryVector) { - return ((VarBinaryVector) vec).get(row); + /** Select the physical vector type once for the entire column. */ + private void readColumn(FieldVector vec, Object[] output) { + if (vec instanceof DateDayVector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = v.get(row); + } + } + } else if (vec instanceof DateMilliVector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = LocalDate.ofEpochDay(Math.floorDiv(v.get(row), 86_400_000L)); + } + } + } else if (vec instanceof TimeStampVector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = v.get(row); + } + } + } else if (vec instanceof BigIntVector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = v.get(row); + } + } + } else if (vec instanceof IntVector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = v.get(row); + } + } + } else if (vec instanceof Float4Vector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = v.get(row); + } + } + } else if (vec instanceof Float8Vector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = v.get(row); + } + } + } else if (vec instanceof BitVector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = v.get(row) != 0; + } + } + } else if (vec instanceof VarCharVector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = new String(v.get(row), StandardCharsets.UTF_8); + } + } + } else if (vec instanceof VarBinaryVector v) { + for (int row = 0; row < output.length; row++) { + if (!v.isNull(row)) { + output[row] = v.get(row); + } + } } else { - Object obj = vec.getObject(row); - return obj != null ? obj.toString() : null; + for (int row = 0; row < output.length; row++) { + if (!vec.isNull(row)) { + Object value = vec.getObject(row); + output[row] = value == null ? null : value.toString(); + } + } } } diff --git a/java/tools/src/main/java/org/apache/tsfile/tools/ParquetSourceReader.java b/java/tools/src/main/java/org/apache/tsfile/tools/ParquetSourceReader.java index 207ef1b60..afe53131c 100644 --- a/java/tools/src/main/java/org/apache/tsfile/tools/ParquetSourceReader.java +++ b/java/tools/src/main/java/org/apache/tsfile/tools/ParquetSourceReader.java @@ -46,6 +46,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.function.Function; public class ParquetSourceReader implements SourceReader { @@ -164,34 +165,29 @@ public SourceBatch readBatch() { Map parquetColIndex = buildParquetColumnIndex(); int numCols = schemaColumnNames.size(); - List rows = new ArrayList<>((int) rowCount); - - for (long r = 0; r < rowCount; r++) { + int count = Math.toIntExact(rowCount); + Object[][] columns = new Object[numCols][count]; + List> extractors = new ArrayList<>(numCols); + for (String name : schemaColumnNames) { + Integer index = parquetColIndex.get(name); + extractors.add(index == null ? null : valueExtractor(index)); + } + // The record reader is row-oriented, but its output can go directly into SourceBatch columns. + for (int row = 0; row < count; row++) { Group group = recordReader.read(); - Object[] row = new Object[numCols]; - - for (int c = 0; c < numCols; c++) { - String colName = schemaColumnNames.get(c); - Integer pIdx = parquetColIndex.get(colName); - if (pIdx == null) { - row[c] = null; - continue; - } - - try { - if (group.getFieldRepetitionCount(pIdx) == 0) { - row[c] = null; - } else { - row[c] = extractValue(group, pIdx); + for (int col = 0; col < numCols; col++) { + Function extractor = extractors.get(col); + if (extractor != null) { + try { + columns[col][row] = extractor.apply(group); + } catch (RuntimeException e) { + // Preserve the existing treatment of missing or malformed field values as null. + columns[col][row] = null; } - } catch (RuntimeException e) { - row[c] = null; } } - rows.add(row); } - - return SourceBatch.fromRows(schemaColumnNames, rows); + return new SourceBatch(schemaColumnNames.toArray(new String[0]), columns, count); } catch (IOException e) { LOGGER.error( Messages.format("log.tools.parquet_read_error", sourceFile.getAbsolutePath()), e); @@ -281,37 +277,30 @@ private Map buildParquetColumnIndex() { return index; } - private Object extractValue(Group group, int fieldIndex) { + private Function valueExtractor(int fieldIndex) { Type fieldType = parquetSchema.getType(fieldIndex); + Function extractor; if (!fieldType.isPrimitive()) { - return group.getGroup(fieldIndex, 0).toString(); - } - - PrimitiveType pt = fieldType.asPrimitiveType(); - switch (pt.getPrimitiveTypeName()) { - case BOOLEAN: - return group.getBoolean(fieldIndex, 0); - case INT32: - return group.getInteger(fieldIndex, 0); - case INT64: - return group.getLong(fieldIndex, 0); - case FLOAT: - return group.getFloat(fieldIndex, 0); - case DOUBLE: - return group.getDouble(fieldIndex, 0); - case BINARY: - case FIXED_LEN_BYTE_ARRAY: - LogicalTypeAnnotation logical = pt.getLogicalTypeAnnotation(); - if (logical instanceof LogicalTypeAnnotation.StringLogicalTypeAnnotation) { - return group.getString(fieldIndex, 0); - } - return group.getBinary(fieldIndex, 0).getBytes(); - case INT96: - // Use getInt96 — INT96 values are Int96Value, not BinaryValue, so getBinary throws CCE. - return int96ToEpochNanos(group.getInt96(fieldIndex, 0).getBytes()); - default: - return group.getValueToString(fieldIndex, 0); + extractor = group -> group.getGroup(fieldIndex, 0).toString(); + } else { + PrimitiveType pt = fieldType.asPrimitiveType(); + extractor = + switch (pt.getPrimitiveTypeName()) { + case BOOLEAN -> group -> group.getBoolean(fieldIndex, 0); + case INT32 -> group -> group.getInteger(fieldIndex, 0); + case INT64 -> group -> group.getLong(fieldIndex, 0); + case FLOAT -> group -> group.getFloat(fieldIndex, 0); + case DOUBLE -> group -> group.getDouble(fieldIndex, 0); + case BINARY, FIXED_LEN_BYTE_ARRAY -> + pt.getLogicalTypeAnnotation() + instanceof LogicalTypeAnnotation.StringLogicalTypeAnnotation + ? group -> group.getString(fieldIndex, 0) + : group -> group.getBinary(fieldIndex, 0).getBytes(); + // INT96 values require getInt96 rather than getBinary. + case INT96 -> group -> int96ToEpochNanos(group.getInt96(fieldIndex, 0).getBytes()); + }; } + return group -> group.getFieldRepetitionCount(fieldIndex) == 0 ? null : extractor.apply(group); } /** diff --git a/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java b/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java index e55019627..43ff5ece1 100644 --- a/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java +++ b/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java @@ -22,6 +22,8 @@ import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.file.metadata.TableSchema; import org.apache.tsfile.i18n.Messages; +import org.apache.tsfile.read.common.type.Type; +import org.apache.tsfile.utils.BitMap; import org.apache.tsfile.write.record.Tablet; import org.apache.tsfile.write.schema.IMeasurementSchema; import org.apache.tsfile.write.schema.MeasurementSchema; @@ -78,11 +80,18 @@ public Tablet build(SourceBatch batch) { for (int col = 0; col < tableSchema.getColumnSchemas().size(); col++) { IMeasurementSchema colSchema = tableSchema.getColumnSchemas().get(col); String colName = colSchema.getMeasurementName(); + Type type = Type.fromTsDataType(colSchema.getType()); + Object targetValues = tablet.getValues()[col]; + // addTimestamp initialized the API bitmap; resolve it and the column type only once. + BitMap nulls = rowCount == 0 ? null : tablet.getBitMaps()[col]; if (tagDefaults.containsKey(colName)) { Object defaultValue = tagDefaults.get(colName); for (int i = 0; i < rowCount; i++) { - tablet.addValue(colName, i, defaultValue); + type.addValue(i, defaultValue, targetValues); + if (defaultValue != null) { + nulls.unmark(i); + } } continue; } diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java index 6b66792ef..00e97cb71 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java @@ -207,6 +207,9 @@ public int write(Tablet tablet, int startRowIndex, int endRowIndex) emptyValueChunkWriters.add(entry.getValue()); } } + // Resolve lazily on the first row to preserve timestamp validation and late-column catch-up. + ValueChunkWriter[] writers = new ValueChunkWriter[measurementSchemas.size()]; + Type[] types = new Type[measurementSchemas.size()]; // TODO: changing to a column-first style by calculating the remaining page space of each // column firsts for (int row = startRowIndex; row < endRowIndex; row++) { @@ -223,10 +226,14 @@ public int write(Tablet tablet, int startRowIndex, int endRowIndex) && tablet.getBitMaps()[columnIndex] != null && tablet.getBitMaps()[columnIndex].isMarked(row); // check isNull by bitMap in tablet - ValueChunkWriter valueChunkWriter = - tryToAddSeriesWriterInternal(measurementSchemas.get(columnIndex)); - Type.fromTsDataType(measurementSchemas.get(columnIndex).getType()) - .write(valueChunkWriter, time, tablet.getValues()[columnIndex], row, isNull); + ValueChunkWriter valueChunkWriter = writers[columnIndex]; + if (valueChunkWriter == null) { + valueChunkWriter = tryToAddSeriesWriterInternal(measurementSchemas.get(columnIndex)); + writers[columnIndex] = valueChunkWriter; + types[columnIndex] = Type.fromTsDataType(measurementSchemas.get(columnIndex).getType()); + } + types[columnIndex].write( + valueChunkWriter, time, tablet.getValues()[columnIndex], row, isNull); } // TODO: we can write the null columns after whole insertion, according to the point number // in the time chunk before and after, no need to do it in a row-by-row manner diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/NonAlignedChunkGroupWriterImpl.java b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/NonAlignedChunkGroupWriterImpl.java index ea73884d2..64f4b6297 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/NonAlignedChunkGroupWriterImpl.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/NonAlignedChunkGroupWriterImpl.java @@ -26,6 +26,8 @@ import org.apache.tsfile.file.metadata.IDeviceID; import org.apache.tsfile.i18n.Messages; import org.apache.tsfile.read.common.type.Type; +import org.apache.tsfile.utils.BitMap; +import org.apache.tsfile.utils.TypeServices; import org.apache.tsfile.write.record.Tablet; import org.apache.tsfile.write.record.datapoint.DataPoint; import org.apache.tsfile.write.schema.IMeasurementSchema; @@ -116,20 +118,35 @@ public int write(Tablet tablet, int startRowIndex, int endRowIndex) String measurementId = timeseries.get(column).getMeasurementName(); Type type = Type.fromTsDataType(timeseries.get(column).getType()); ChunkWriterImpl chunkWriter = chunkWriters.get(measurementId); - pointCount = 0; - for (int row = startRowIndex; row < endRowIndex; row++) { - // check isNull in tablet - if (tablet.getBitMaps() != null - && tablet.getBitMaps()[column] != null - && tablet.getBitMaps()[column].isMarked(row)) { - continue; + BitMap nulls = tablet.getBitMaps() == null ? null : tablet.getBitMaps()[column]; + int firstRow = startRowIndex; + while (firstRow < endRowIndex && nulls != null && nulls.isMarked(firstRow)) { + firstRow++; + } + // An empty column must not inspect its value array or invoke a value writer. + if (firstRow >= endRowIndex) { + continue; + } + TabletWriteContext context = + new TabletWriteContext(deviceId, measurementId, lastTimeMap.get(measurementId)); + try { + TypeServices.WRITE_TABLET_COLUMN_SERVICE + .call(type) + .write( + chunkWriter, + tablet.getTimestamps(), + tablet.getValues()[column], + nulls, + firstRow, + endRowIndex, + context); + } finally { + // Commit the successfully written prefix even when a later row is out of order. + if (context.getPointCount() != 0) { + lastTimeMap.put(measurementId, context.getLastTime()); } - long time = tablet.getTimestamps()[row]; - checkIsHistoryData(measurementId, time); - pointCount++; - type.write(chunkWriter, time, tablet.getValues()[column], row); - lastTimeMap.put(measurementId, time); } + pointCount = context.getPointCount(); maxPointCount = Math.max(pointCount, maxPointCount); } return maxPointCount; diff --git a/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/AlignedNullWriteTest.java b/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/AlignedNullWriteTest.java index 8e6eebe36..4448ae31e 100644 --- a/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/AlignedNullWriteTest.java +++ b/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/AlignedNullWriteTest.java @@ -22,20 +22,77 @@ import org.apache.tsfile.common.conf.TSFileDescriptor; import org.apache.tsfile.encoding.decoder.PlainDecoder; import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.exception.write.WriteProcessException; import org.apache.tsfile.file.metadata.IDeviceID; import org.apache.tsfile.file.metadata.enums.CompressionType; import org.apache.tsfile.file.metadata.enums.TSEncoding; +import org.apache.tsfile.write.record.Tablet; import org.apache.tsfile.write.record.datapoint.IntDataPoint; import org.apache.tsfile.write.schema.MeasurementSchema; import org.junit.Test; import java.nio.ByteBuffer; +import java.util.Arrays; import java.util.Collections; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThrows; public class AlignedNullWriteTest { + @Test + public void tabletWriterCacheKeepsLateColumnPagesAndFailureOrder() throws Exception { + TSFileConfig config = TSFileDescriptor.getInstance().getConfig(); + int oldLimit = config.getMaxNumberOfPointsInPage(); + try { + config.setMaxNumberOfPointsInPage(8); + AlignedChunkGroupWriterImpl scalar = + new AlignedChunkGroupWriterImpl(IDeviceID.Factory.DEFAULT_FACTORY.create("root.tablet")); + AlignedChunkGroupWriterImpl batch = + new AlignedChunkGroupWriterImpl(IDeviceID.Factory.DEFAULT_FACTORY.create("root.tablet")); + MeasurementSchema missing = + new MeasurementSchema( + "missing", TSDataType.INT32, TSEncoding.PLAIN, CompressionType.UNCOMPRESSED); + scalar.tryToAddSeriesWriter(missing); + batch.tryToAddSeriesWriter(missing); + for (int row = 0; row < 11; row++) { + scalar.write(row, Collections.emptyList()); + batch.write(row, Collections.emptyList()); + } + MeasurementSchema late = + new MeasurementSchema( + "late", TSDataType.INT32, TSEncoding.PLAIN, CompressionType.UNCOMPRESSED); + scalar.tryToAddSeriesWriter(late); + Tablet tablet = new Tablet("root.tablet", Collections.singletonList(late), 24); + for (int row = 0; row < 24; row++) { + tablet.addTimestamp(row, row + 11); + tablet.addValue("late", row, row); + scalar.write(row + 11, Collections.singletonList(new IntDataPoint("late", row))); + } + batch.write(tablet); + for (String name : Arrays.asList("missing", "late")) { + ValueChunkWriter expected = scalar.valueChunkWriterMap.get(name); + ValueChunkWriter actual = batch.valueChunkWriterMap.get(name); + expected.sealCurrentPage(); + actual.sealCurrentPage(); + assertEquals(expected.getByteBuffer(), actual.getByteBuffer()); + assertEquals(expected.getNumOfPages(), actual.getNumOfPages()); + assertEquals(expected.getStatistics(), actual.getStatistics()); + } + Tablet rejected = + new Tablet( + "root.tablet", + Collections.singletonList(new MeasurementSchema("rejected", TSDataType.INT32)), + 1); + rejected.addTimestamp(0, 0); + rejected.addValue("rejected", 0, 42); + assertThrows(WriteProcessException.class, () -> batch.write(rejected)); + assertEquals(2, batch.valueChunkWriterMap.size()); + } finally { + config.setMaxNumberOfPointsInPage(oldLimit); + } + } + @Test public void lateColumnsAndMissingRowsKeepPageBoundaries() throws Exception { TSFileConfig config = TSFileDescriptor.getInstance().getConfig(); From bf46f2c52e97b7e34c5b8dbb592bce64ef399b20 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Wed, 16 Sep 2026 18:15:38 +0800 Subject: [PATCH 3/4] tmp --- .../apache/tsfile/tools/TabletBuilder.java | 5 +- .../tsfile/tools/TabletBuilderTest.java | 73 ++++++++ .../templates/TabletColumnWriterTemplate.ftl | 72 ++++++++ .../org/apache/tsfile/utils/TypeServices.java | 39 +++++ .../write/chunk/TabletWriteContext.java | 68 ++++++++ .../write/chunk/TabletColumnWriteTest.java | 164 ++++++++++++++++++ 6 files changed, 420 insertions(+), 1 deletion(-) create mode 100644 java/tsfile/src/main/codegen/templates/TabletColumnWriterTemplate.ftl create mode 100644 java/tsfile/src/main/java/org/apache/tsfile/write/chunk/TabletWriteContext.java create mode 100644 java/tsfile/src/test/java/org/apache/tsfile/write/chunk/TabletColumnWriteTest.java diff --git a/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java b/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java index 43ff5ece1..61d4457d7 100644 --- a/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java +++ b/java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java @@ -117,7 +117,10 @@ public Tablet build(SourceBatch batch) { colSchema.getType(), isMeasurement, importSchema.getTimePrecision()); } Object converted = converter.apply(rawValue); - tablet.addValue(colName, i, converted); + type.addValue(i, converted, targetValues); + if (converted != null) { + nulls.unmark(i); + } } } diff --git a/java/tools/src/test/java/org/apache/tsfile/tools/TabletBuilderTest.java b/java/tools/src/test/java/org/apache/tsfile/tools/TabletBuilderTest.java index 5dcfbb3ef..3e1c1ae00 100644 --- a/java/tools/src/test/java/org/apache/tsfile/tools/TabletBuilderTest.java +++ b/java/tools/src/test/java/org/apache/tsfile/tools/TabletBuilderTest.java @@ -21,6 +21,7 @@ import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.write.record.Tablet; +import org.apache.tsfile.write.schema.IMeasurementSchema; import org.junit.Test; @@ -35,6 +36,78 @@ public class TabletBuilderTest { + @Test + public void testColumnInsertionMatchesPublicTabletApi() { + TSDataType[] types = { + TSDataType.BOOLEAN, + TSDataType.INT32, + TSDataType.INT64, + TSDataType.FLOAT, + TSDataType.DOUBLE, + TSDataType.DATE, + TSDataType.TIMESTAMP, + TSDataType.TEXT, + TSDataType.STRING, + TSDataType.BLOB, + TSDataType.OBJECT + }; + List columns = new ArrayList<>(); + columns.add(new ImportSchema.SourceColumn("time", TSDataType.INT64)); + String[] names = new String[types.length + 1]; + names[0] = "time"; + for (int col = 0; col < types.length; col++) { + names[col + 1] = "s" + col; + columns.add(new ImportSchema.SourceColumn(names[col + 1], types[col])); + } + ImportSchema schema = + buildSchema( + "test", + "time", + Collections.singletonList(new ImportSchema.TagColumn("tag", "default")), + columns.toArray(new ImportSchema.SourceColumn[0])); + schema.setNullFormat("NULL"); + TabletBuilder builder = new TabletBuilder(schema, new TimeConverter("ms")); + for (int count : new int[] {0, 1, 35}) { + Object[][] values = new Object[names.length][count]; + for (int row = 0; row < count; row++) { + values[0][row] = (long) (count - row); + for (int col = 0; col < types.length; col++) { + values[col + 1][row] = + row % 4 == 0 + ? "NULL" + : switch (types[col]) { + case BOOLEAN -> "true"; + case DATE -> "2024-02-29"; + case INT32, INT64, TIMESTAMP, FLOAT, DOUBLE -> row; + case TEXT, STRING, BLOB, OBJECT -> "v" + row; + case VECTOR, UNKNOWN -> throw new AssertionError(); + }; + } + } + Tablet actual = builder.build(new SourceBatch(names, values, count)); + List schemas = builder.getTableSchema().getColumnSchemas(); + Tablet expected = + new Tablet( + "test", + IMeasurementSchema.getMeasurementNameList(schemas), + IMeasurementSchema.getDataTypeList(schemas), + builder.getTableSchema().getColumnTypes(), + count); + for (int row = 0; row < count; row++) { + expected.addTimestamp(row, row + 1); + expected.addValue("tag", row, "default"); + for (int col = 0; col < types.length; col++) { + Object raw = values[col + 1][count - row - 1]; + if (!"NULL".equals(raw)) { + expected.addValue(names[col + 1], row, ValueConverter.convert(raw, types[col], true)); + } + } + } + expected.setRowSize(count); + assertEquals(expected, actual); + } + } + @Test public void testColumnConversionKeepsSortedRowsAndNulls() { ImportSchema schema = diff --git a/java/tsfile/src/main/codegen/templates/TabletColumnWriterTemplate.ftl b/java/tsfile/src/main/codegen/templates/TabletColumnWriterTemplate.ftl new file mode 100644 index 000000000..ae6a3c083 --- /dev/null +++ b/java/tsfile/src/main/codegen/templates/TabletColumnWriterTemplate.ftl @@ -0,0 +1,72 @@ +<@pp.dropOutputFile /> +<#list [["Boolean", "boolean"], ["Int", "int"], ["Long", "long"], + ["Float", "float"], ["Double", "double"], ["Binary", "Binary"], + ["Date", "LocalDate"]] as t> +<#assign name=t[0] primitive=t[1]> +<@pp.changeOutputFile name="/org/apache/tsfile/read/common/type/service/${name}TabletColumnWriter.java" /> +/* + * 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.tsfile.read.common.type.service; + +import org.apache.tsfile.exception.write.WriteProcessException; +import org.apache.tsfile.utils.BitMap; +<#if name == "Binary"> +import org.apache.tsfile.utils.Binary; + +<#if name == "Date"> +import org.apache.tsfile.utils.DateUtils; +import java.time.LocalDate; + +import org.apache.tsfile.utils.TypeServices.TabletColumnWriter; +import org.apache.tsfile.write.chunk.ChunkWriterImpl; +import org.apache.tsfile.write.chunk.TabletWriteContext; + +/** Generated typed loop; retain scalar encoding to preserve page boundaries and SDT state. */ +public final class ${name}TabletColumnWriter implements TabletColumnWriter { + public static final ${name}TabletColumnWriter INSTANCE = new ${name}TabletColumnWriter(); + + private ${name}TabletColumnWriter() {} + + @Override + public void write(ChunkWriterImpl writer, long[] times, Object values, BitMap nulls, + int start, int end, TabletWriteContext context) throws WriteProcessException { +<#if name == "Date"> + if (values instanceof int[]) { + IntTabletColumnWriter.INSTANCE.write(writer, times, values, nulls, start, end, context); + return; + } + + ${primitive}[] column = (${primitive}[]) values; + for (int row = start; row < end; row++) { + if (nulls != null && nulls.isMarked(row)) { + continue; + } + long time = times[row]; + context.checkTime(time); +<#if name == "Date"> + writer.write(time, DateUtils.parseDateExpressionToInt(column[row])); +<#else> + writer.write(time, column[row]); + + context.recordTime(time); + } + } +} + diff --git a/java/tsfile/src/main/java/org/apache/tsfile/utils/TypeServices.java b/java/tsfile/src/main/java/org/apache/tsfile/utils/TypeServices.java index 01763e1af..0afe8b77a 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/utils/TypeServices.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/utils/TypeServices.java @@ -22,21 +22,31 @@ import org.apache.tsfile.block.column.ColumnBuilder; import org.apache.tsfile.encoding.decoder.Decoder; import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.exception.write.WriteProcessException; import org.apache.tsfile.i18n.Messages; import org.apache.tsfile.read.common.BatchData; import org.apache.tsfile.read.common.block.TsBlockBuilder; import org.apache.tsfile.read.common.block.column.BinaryColumnBuilder; import org.apache.tsfile.read.common.type.service.BinaryPageDataReader; +import org.apache.tsfile.read.common.type.service.BinaryTabletColumnWriter; import org.apache.tsfile.read.common.type.service.BooleanPageDataReader; +import org.apache.tsfile.read.common.type.service.BooleanTabletColumnWriter; import org.apache.tsfile.read.common.type.service.DatePageDataReader; +import org.apache.tsfile.read.common.type.service.DateTabletColumnWriter; import org.apache.tsfile.read.common.type.service.DoublePageDataReader; +import org.apache.tsfile.read.common.type.service.DoubleTabletColumnWriter; import org.apache.tsfile.read.common.type.service.FloatPageDataReader; +import org.apache.tsfile.read.common.type.service.FloatTabletColumnWriter; import org.apache.tsfile.read.common.type.service.IntPageDataReader; +import org.apache.tsfile.read.common.type.service.IntTabletColumnWriter; import org.apache.tsfile.read.common.type.service.LongPageDataReader; +import org.apache.tsfile.read.common.type.service.LongTabletColumnWriter; import org.apache.tsfile.read.common.type.service.TypeService; import org.apache.tsfile.read.filter.basic.Filter; import org.apache.tsfile.read.reader.series.PaginationController; import org.apache.tsfile.write.UnSupportedDataTypeException; +import org.apache.tsfile.write.chunk.ChunkWriterImpl; +import org.apache.tsfile.write.chunk.TabletWriteContext; import org.apache.tsfile.write.chunk.ValueChunkWriter; import java.io.IOException; @@ -61,7 +71,36 @@ public final class TypeServices { .setChecked(true); }; + /** Resolve the Tablet column type before entering its row loop. */ + public static final TypeService WRITE_TABLET_COLUMN_SERVICE = + type -> + switch (type.getTypeEnum()) { + case BOOLEAN -> BooleanTabletColumnWriter.INSTANCE; + case INT32 -> IntTabletColumnWriter.INSTANCE; + case DATE -> DateTabletColumnWriter.INSTANCE; + case INT64, TIMESTAMP -> LongTabletColumnWriter.INSTANCE; + case FLOAT -> FloatTabletColumnWriter.INSTANCE; + case DOUBLE -> DoubleTabletColumnWriter.INSTANCE; + case TEXT, BLOB, STRING, OBJECT -> BinaryTabletColumnWriter.INSTANCE; + case ROW, UNKNOWN, VECTOR -> + throw new UnSupportedDataTypeException(String.valueOf(type.getTypeEnum())) + .setChecked(true); + }; + + public interface TabletColumnWriter { + void write( + ChunkWriterImpl writer, + long[] times, + Object values, + BitMap nulls, + int start, + int end, + TabletWriteContext context) + throws WriteProcessException; + } + static { + WRITE_TABLET_COLUMN_SERVICE.check(); READ_PAGE_BATCH_SERVICE.check(); } diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/TabletWriteContext.java b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/TabletWriteContext.java new file mode 100644 index 000000000..3a97619dc --- /dev/null +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/TabletWriteContext.java @@ -0,0 +1,68 @@ +/* + * 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.tsfile.write.chunk; + +import org.apache.tsfile.common.constant.TsFileConstant; +import org.apache.tsfile.exception.write.WriteProcessException; +import org.apache.tsfile.file.metadata.IDeviceID; +import org.apache.tsfile.i18n.Messages; + +/** Per-column state shared by the type-specialized Tablet write loops. */ +public final class TabletWriteContext { + private final IDeviceID deviceId; + private final String measurementId; + private boolean hasLastTime; + private long lastTime; + private int pointCount; + + public TabletWriteContext(IDeviceID deviceId, String measurementId, Long lastTime) { + this.deviceId = deviceId; + this.measurementId = measurementId; + this.hasLastTime = lastTime != null; + this.lastTime = lastTime == null ? 0 : lastTime; + } + + public void checkTime(long time) throws WriteProcessException { + if (hasLastTime && time <= lastTime) { + throw new WriteProcessException( + Messages.format( + "error.write.chunk_group_non_aligned_out_of_order", + deviceId, + TsFileConstant.PATH_SEPARATOR, + measurementId, + lastTime)); + } + } + + /** Advance only after a successful write, so failures retain the last accepted timestamp. */ + public void recordTime(long time) { + lastTime = time; + hasLastTime = true; + pointCount++; + } + + public long getLastTime() { + return lastTime; + } + + public int getPointCount() { + return pointCount; + } +} diff --git a/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/TabletColumnWriteTest.java b/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/TabletColumnWriteTest.java new file mode 100644 index 000000000..af65a99c8 --- /dev/null +++ b/java/tsfile/src/test/java/org/apache/tsfile/write/chunk/TabletColumnWriteTest.java @@ -0,0 +1,164 @@ +/* + * 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.tsfile.write.chunk; + +import org.apache.tsfile.common.conf.TSFileDescriptor; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.exception.write.WriteProcessException; +import org.apache.tsfile.file.metadata.IDeviceID; +import org.apache.tsfile.file.metadata.enums.CompressionType; +import org.apache.tsfile.file.metadata.enums.TSEncoding; +import org.apache.tsfile.read.common.type.Type; +import org.apache.tsfile.utils.Binary; +import org.apache.tsfile.utils.BitMap; +import org.apache.tsfile.utils.TypeServices; +import org.apache.tsfile.write.record.Tablet; +import org.apache.tsfile.write.schema.MeasurementSchema; + +import org.junit.Test; + +import java.nio.charset.StandardCharsets; +import java.time.LocalDate; +import java.util.Collections; +import java.util.Map; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThrows; + +public class TabletColumnWriteTest { + @Test + public void testEncodedPagesMatchScalarForAllTypes() throws Exception { + int oldLimit = TSFileDescriptor.getInstance().getConfig().getMaxNumberOfPointsInPage(); + try { + // Force multiple pages and exercise both bitmap boundaries and a nonzero source offset. + TSFileDescriptor.getInstance().getConfig().setMaxNumberOfPointsInPage(7); + for (TSDataType dataType : TSDataType.values()) { + if (dataType == TSDataType.UNKNOWN || dataType == TSDataType.VECTOR) { + continue; + } + Type type = Type.fromTsDataType(dataType); + for (boolean sparse : new boolean[] {false, true}) { + Object values = type.createArray(35); + long[] times = new long[35]; + BitMap nulls = sparse ? new BitMap(35) : null; + for (int row = 0; row < 35; row++) { + times[row] = row - 20; + type.addValue(row, value(dataType, row), values); + if (sparse && row % 3 == 0) { + nulls.mark(row); + } + } + compare(type, dataType, times, values, nulls); + if (dataType == TSDataType.DATE) { + int[] intDates = new int[35]; + for (int row = 0; row < intDates.length; row++) { + intDates[row] = 20240101 + row % 28; + } + compare(type, dataType, times, intDates, nulls); + } + } + } + } finally { + TSFileDescriptor.getInstance().getConfig().setMaxNumberOfPointsInPage(oldLimit); + } + } + + private void compare(Type type, TSDataType dataType, long[] times, Object values, BitMap nulls) + throws Exception { + MeasurementSchema schema = + new MeasurementSchema("s", dataType, TSEncoding.PLAIN, CompressionType.UNCOMPRESSED); + compareWriters(type, dataType, times, values, nulls, schema); + if (dataType == TSDataType.INT32 + || dataType == TSDataType.INT64 + || dataType == TSDataType.FLOAT + || dataType == TSDataType.DOUBLE) { + schema.setProps(Map.of("loss", "sdt", "compdev", "0.1", "compmaxtime", "5")); + compareWriters(type, dataType, times, values, nulls, schema); + } + } + + private void compareWriters( + Type type, + TSDataType dataType, + long[] times, + Object values, + BitMap nulls, + MeasurementSchema schema) + throws Exception { + ChunkWriterImpl scalar = new ChunkWriterImpl(schema); + ChunkWriterImpl batch = new ChunkWriterImpl(schema); + int count = 0; + for (int row = 2; row < 34; row++) { + if (nulls == null || !nulls.isMarked(row)) { + type.write(scalar, times[row], values, row); + count++; + } + } + TabletWriteContext context = + new TabletWriteContext(IDeviceID.Factory.DEFAULT_FACTORY.create("root.d"), "s", null); + TypeServices.WRITE_TABLET_COLUMN_SERVICE + .call(type) + .write(batch, times, values, nulls, 2, 34, context); + scalar.sealCurrentPage(); + batch.sealCurrentPage(); + assertEquals(dataType.toString(), scalar.getByteBuffer(), batch.getByteBuffer()); + assertEquals(scalar.getStatistics(), batch.getStatistics()); + assertEquals(scalar.getNumOfPages(), batch.getNumOfPages()); + assertEquals(count, context.getPointCount()); + } + + @Test + public void testOutOfOrderKeepsWrittenPrefixAndSkipsNulls() throws Exception { + NonAlignedChunkGroupWriterImpl writer = + new NonAlignedChunkGroupWriterImpl(IDeviceID.Factory.DEFAULT_FACTORY.create("root.d")); + writer.tryToAddSeriesWriter(new MeasurementSchema("s", TSDataType.INT32)); + Tablet tablet = + new Tablet( + "root.d", Collections.singletonList(new MeasurementSchema("s", TSDataType.INT32)), 5); + long[] times = {Long.MIN_VALUE, 10, 3, 9, 20}; + for (int row = 0; row < times.length; row++) { + tablet.addTimestamp(row, times[row]); + if (row != 2) { + tablet.addValue("s", row, row); + } + } + assertThrows(WriteProcessException.class, () -> writer.write(tablet, 0, 5)); + assertEquals(Long.valueOf(10), writer.getLastTimeMap().get("s")); + assertEquals(1, writer.write(tablet, 4, 5)); + assertEquals(Long.valueOf(20), writer.getLastTimeMap().get("s")); + assertEquals(0, writer.write(tablet, 2, 3)); + assertEquals(Long.valueOf(20), writer.getLastTimeMap().get("s")); + tablet.getValues()[0] = null; + assertEquals(0, writer.write(tablet, 2, 3)); + assertEquals(0, writer.write(tablet, 0, 0)); + } + + private Object value(TSDataType type, int row) { + return switch (type) { + case BOOLEAN -> row % 2 == 0; + case INT32 -> row; + case DATE -> LocalDate.of(2024, 1, 1 + row % 28); + case INT64, TIMESTAMP -> (long) row; + case FLOAT -> row * 1.5f; + case DOUBLE -> row * 1.5; + case TEXT, STRING, BLOB, OBJECT -> new Binary("v" + row, StandardCharsets.UTF_8); + case VECTOR, UNKNOWN -> throw new AssertionError(type); + }; + } +} From 5d15440e831dfeb00f025d8c78c716945cc250c2 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Thu, 17 Sep 2026 10:07:34 +0800 Subject: [PATCH 4/4] tmp --- .../reader/block/SingleDeviceTsBlockReader.java | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/block/SingleDeviceTsBlockReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/block/SingleDeviceTsBlockReader.java index e149a648a..da4f05e11 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/block/SingleDeviceTsBlockReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/block/SingleDeviceTsBlockReader.java @@ -234,8 +234,9 @@ private void fillIdColumn(Column column, Object val, int startPos, int endPos) { column.setPositionCount(endPos); } - private static void fillSingleMeasurementColumn(Column column, BatchData batchData, int pos) { - Type.fromTsDataType(batchData.getDataType()).setTo(batchData, column, pos); + private static void fillSingleMeasurementColumn( + Column column, BatchData batchData, Type type, int pos) { + type.setTo(batchData, column, pos); column.setPositionCount(pos + 1); } @@ -274,6 +275,8 @@ public static class SingleMeasurementColumnContext extends MeasurementColumnCont private final String columnName; private final List posInResult; + // A measurement context keeps one data type across all batches, so bind its Type once. + private final Type type; public SingleMeasurementColumnContext( String columnName, @@ -283,6 +286,7 @@ public SingleMeasurementColumnContext( super(seriesReader, currentBatch); this.columnName = columnName; this.posInResult = posInResult; + this.type = Type.fromTsDataType(currentBatch.getDataType()); } @Override @@ -294,7 +298,7 @@ void removeFrom(Map columnContextMap) { void fillInto(TsBlock block, int position) { for (Integer pos : posInResult) { final Column column = block.getColumn(pos); - fillSingleMeasurementColumn(column, currentBatch, position); + fillSingleMeasurementColumn(column, currentBatch, type, position); } } } @@ -322,10 +326,10 @@ void fillInto(TsBlock block, int blockRowNum) { for (int i = 0; i < vector.length; i++) { final TsPrimitiveType value = vector[i]; final List columnPositions = posInResult.get(i); + final Type type = value == null ? null : Type.fromTsDataType(value.getDataType()); for (Integer pos : columnPositions) { if (value != null) { - Type.fromTsDataType(value.getDataType()) - .setTo(value, block.getColumn(pos), blockRowNum); + type.setTo(value, block.getColumn(pos), blockRowNum); } else { block.getColumn(pos).setNull(blockRowNum, blockRowNum + 1); }