Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -182,23 +182,15 @@ public SourceBatch readBatch() {
}

int numCols = schemaColumnNames.size();
List<Object[]> 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];

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Build SourceBatch columns directly because both Arrow and the destination are columnar. Resolving each vector once removes repeated name lookups and the intermediate row arrays/transpose; absent vectors remain null. Existing Arrow reader tests and the cross-version output checks passed.

// 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;
Expand Down Expand Up @@ -293,40 +285,75 @@ private List<String> 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) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Bind the physical vector type before iterating through its values. Each typed loop retains null checks, native timestamp precision, date conversion, UTF-8 decoding, and the generic fallback, while avoiding an instanceof chain for every cell.

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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Expand Down Expand Up @@ -164,34 +165,29 @@ public SourceBatch readBatch() {
Map<String, Integer> parquetColIndex = buildParquetColumnIndex();

int numCols = schemaColumnNames.size();
List<Object[]> rows = new ArrayList<>((int) rowCount);

for (long r = 0; r < rowCount; r++) {
int count = Math.toIntExact(rowCount);
Object[][] columns = new Object[numCols][count];

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Keep the Parquet record reader row-oriented but write directly to the final column arrays. Column indices and extractors are resolved once per batch; missing or malformed fields retain the existing null behavior. This removes the row transpose without changing the input reader protocol.

List<Function<Group, Object>> 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<Group, Object> 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);
Expand Down Expand Up @@ -281,37 +277,30 @@ private Map<String, Integer> buildParquetColumnIndex() {
return index;
}

private Object extractValue(Group group, int fieldIndex) {
private Function<Group, Object> valueExtractor(int fieldIndex) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Choose physical/logical-type conversion once per column, including the dedicated INT96 accessor, rather than inspecting schema metadata for every value. The wrapper still checks field presence per row. Local benchmarks showed lower allocation but inconclusive elapsed-time changes for Parquet.

Type fieldType = parquetSchema.getType(fieldIndex);
Function<Group, Object> 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);
}

/**
Expand Down
60 changes: 43 additions & 17 deletions java/tools/src/main/java/org/apache/tsfile/tools/TabletBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -30,6 +32,7 @@
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Function;

public class TabletBuilder {

Expand Down Expand Up @@ -70,31 +73,54 @@ 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();

if (tagDefaults.containsKey(colName)) {
tablet.addValue(colName, i, tagDefaults.get(colName));
continue;
// 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();
Type type = Type.fromTsDataType(colSchema.getType());

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Traverse columns after establishing sorted timestamps so schema lookup, the destination array, and type resolution are reused across rows. Direct insertion preserves Tablet bitmap semantics by unmarking only non-null converted values; tag defaults and sorted source indices retain their original behavior. Conversion is initialized lazily so empty/all-null columns do not invoke converters. The new tests compare this path with the public Tablet insertion API.

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++) {
type.addValue(i, defaultValue, targetValues);
if (defaultValue != null) {
nulls.unmark(i);
}
}
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<Object, Object> 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());
tablet.addValue(colName, i, converted);
// 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);
type.addValue(i, converted, targetValues);
if (converted != null) {
nulls.unmark(i);
}
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Expand Down Expand Up @@ -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<Object, Object> converterFor(

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Resolve the string and object conversion services once for a fixed target column. The returned converter still selects the appropriate input representation for each cell, so columns mixing strings and native Java values keep the same conversion rules and timestamp precision.

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
Expand Down
Loading
Loading