This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 9a5661e5f1 [format] Support narrowing Parquet reads for TINYINT and
SMALLINT (#8629)
9a5661e5f1 is described below
commit 9a5661e5f18b7908851308f77c6686a3ade8cc7a
Author: Eunbin Son <[email protected]>
AuthorDate: Wed Jul 15 09:34:00 2026 +0900
[format] Support narrowing Parquet reads for TINYINT and SMALLINT (#8629)
---
.../reader/ParquetVectorUpdaterFactory.java | 86 +++++++++++++++++++++-
.../reader/FileTypeNotMatchReadTypeTest.java | 78 ++++++++++++++++++++
2 files changed, 162 insertions(+), 2 deletions(-)
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetVectorUpdaterFactory.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetVectorUpdaterFactory.java
index afa60ff070..3dcd3fbd03 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetVectorUpdaterFactory.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/ParquetVectorUpdaterFactory.java
@@ -147,12 +147,24 @@ public class ParquetVectorUpdaterFactory {
@Override
public UpdaterFactory visit(TinyIntType tinyIntType) {
- return c -> new ByteUpdater();
+ return c -> {
+ if (c.getPrimitiveType().getPrimitiveTypeName()
+ == PrimitiveType.PrimitiveTypeName.INT64) {
+ return new ByteFromLongUpdater();
+ }
+ return new ByteUpdater();
+ };
}
@Override
public UpdaterFactory visit(SmallIntType smallIntType) {
- return c -> new ShortUpdater();
+ return c -> {
+ if (c.getPrimitiveType().getPrimitiveTypeName()
+ == PrimitiveType.PrimitiveTypeName.INT64) {
+ return new ShortFromLongUpdater();
+ }
+ return new ShortUpdater();
+ };
}
@Override
@@ -417,6 +429,41 @@ public class ParquetVectorUpdaterFactory {
}
}
+ private static class ByteFromLongUpdater implements
ParquetVectorUpdater<WritableByteVector> {
+ @Override
+ public void readValues(
+ int total,
+ int offset,
+ WritableByteVector values,
+ VectorizedValuesReader valuesReader) {
+ for (int i = 0; i < total; i++) {
+ values.setByte(offset + i, (byte)
Math.toIntExact(valuesReader.readLong()));
+ }
+ }
+
+ @Override
+ public void skipValues(int total, VectorizedValuesReader valuesReader)
{
+ valuesReader.skipLongs(total);
+ }
+
+ @Override
+ public void readValue(
+ int offset, WritableByteVector values, VectorizedValuesReader
valuesReader) {
+ values.setByte(offset, (byte)
Math.toIntExact(valuesReader.readLong()));
+ }
+
+ @Override
+ public void decodeSingleDictionaryId(
+ int offset,
+ WritableByteVector values,
+ WritableIntVector dictionaryIds,
+ Dictionary dictionary) {
+ values.setByte(
+ offset,
+ (byte)
Math.toIntExact(dictionary.decodeToLong(dictionaryIds.getInt(offset))));
+ }
+ }
+
private static class ShortUpdater implements
ParquetVectorUpdater<WritableShortVector> {
@Override
public void readValues(
@@ -450,6 +497,41 @@ public class ParquetVectorUpdaterFactory {
}
}
+ private static class ShortFromLongUpdater implements
ParquetVectorUpdater<WritableShortVector> {
+ @Override
+ public void readValues(
+ int total,
+ int offset,
+ WritableShortVector values,
+ VectorizedValuesReader valuesReader) {
+ for (int i = 0; i < total; i++) {
+ values.setShort(offset + i, (short)
Math.toIntExact(valuesReader.readLong()));
+ }
+ }
+
+ @Override
+ public void skipValues(int total, VectorizedValuesReader valuesReader)
{
+ valuesReader.skipLongs(total);
+ }
+
+ @Override
+ public void readValue(
+ int offset, WritableShortVector values, VectorizedValuesReader
valuesReader) {
+ values.setShort(offset, (short)
Math.toIntExact(valuesReader.readLong()));
+ }
+
+ @Override
+ public void decodeSingleDictionaryId(
+ int offset,
+ WritableShortVector values,
+ WritableIntVector dictionaryIds,
+ Dictionary dictionary) {
+ values.setShort(
+ offset,
+ (short)
Math.toIntExact(dictionary.decodeToLong(dictionaryIds.getInt(offset))));
+ }
+ }
+
private static class LongUpdater implements
ParquetVectorUpdater<WritableLongVector> {
@Override
public void readValues(
diff --git
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/reader/FileTypeNotMatchReadTypeTest.java
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/reader/FileTypeNotMatchReadTypeTest.java
index 964f4a21eb..479fcd2038 100644
---
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/reader/FileTypeNotMatchReadTypeTest.java
+++
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/reader/FileTypeNotMatchReadTypeTest.java
@@ -237,6 +237,84 @@ public class FileTypeNotMatchReadTypeTest {
file.delete();
}
+ @Test
+ public void testReadByteFromInt64() throws Exception {
+ String fileName = "test.parquet";
+ String fileWholePath = tempDir + "/" + fileName;
+
+ RowType rowTypeWrite = RowType.of(new DataField(0, "byte_col",
DataTypes.BIGINT()));
+ RowType rowTypeRead = RowType.of(new DataField(0, "byte_col",
DataTypes.TINYINT()));
+ MessageType messageType =
Util.convertToParquetMessageType(rowTypeWrite);
+ ParquetRowDataBuilderForTest parquetRowDataBuilder =
+ new ParquetRowDataBuilderForTest(
+ new LocalOutputFile(new
File(fileWholePath).toPath()),
+ rowTypeWrite,
+ messageType)
+ .enableDictionaryEncoding();
+ ParquetWriter<InternalRow> parquetWriter =
parquetRowDataBuilder.build();
+
+ for (int i = 0; i < 100; i++) {
+ parquetWriter.write(GenericRow.of((long) i));
+ }
+ parquetWriter.close();
+
+ ParquetReaderFactory parquetReaderFactory =
+ new ParquetReaderFactory(new Options(), rowTypeRead, 100,
null);
+
+ File file = new File(fileWholePath);
+ FileRecordReader<InternalRow> fileRecordReader =
+ parquetReaderFactory.createReader(
+ new FormatReaderContext(
+ LocalFileIO.create(),
+ new
org.apache.paimon.fs.Path(tempDir.toString(), fileName),
+ file.length()));
+
+ FileRecordIterator<InternalRow> batch = fileRecordReader.readBatch();
+ for (int i = 0; i < 100; i++) {
+ assertThat(batch.next().getByte(0)).isEqualTo((byte) i);
+ }
+ file.delete();
+ }
+
+ @Test
+ public void testReadShortFromInt64() throws Exception {
+ String fileName = "test.parquet";
+ String fileWholePath = tempDir + "/" + fileName;
+
+ RowType rowTypeWrite = RowType.of(new DataField(0, "short_col",
DataTypes.BIGINT()));
+ RowType rowTypeRead = RowType.of(new DataField(0, "short_col",
DataTypes.SMALLINT()));
+ MessageType messageType =
Util.convertToParquetMessageType(rowTypeWrite);
+ ParquetRowDataBuilderForTest parquetRowDataBuilder =
+ new ParquetRowDataBuilderForTest(
+ new LocalOutputFile(new
File(fileWholePath).toPath()),
+ rowTypeWrite,
+ messageType)
+ .enableDictionaryEncoding();
+ ParquetWriter<InternalRow> parquetWriter =
parquetRowDataBuilder.build();
+
+ for (int i = 0; i < 100; i++) {
+ parquetWriter.write(GenericRow.of((long) i));
+ }
+ parquetWriter.close();
+
+ ParquetReaderFactory parquetReaderFactory =
+ new ParquetReaderFactory(new Options(), rowTypeRead, 100,
null);
+
+ File file = new File(fileWholePath);
+ FileRecordReader<InternalRow> fileRecordReader =
+ parquetReaderFactory.createReader(
+ new FormatReaderContext(
+ LocalFileIO.create(),
+ new
org.apache.paimon.fs.Path(tempDir.toString(), fileName),
+ file.length()));
+
+ FileRecordIterator<InternalRow> batch = fileRecordReader.readBatch();
+ for (int i = 0; i < 100; i++) {
+ assertThat(batch.next().getShort(0)).isEqualTo((short) i);
+ }
+ file.delete();
+ }
+
@Test
public void testArray() throws Exception {
String fileName = "test.parquet";