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";

Reply via email to