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 5722fbc85f [format] Avoid lossy Parquet predicate pushdown (#8999)
5722fbc85f is described below

commit 5722fbc85f9085cc3f309db2bc4e2c004901aa2a
Author: Jingsong Lee <[email protected]>
AuthorDate: Tue Aug 4 08:43:42 2026 +0800

    [format] Avoid lossy Parquet predicate pushdown (#8999)
---
 .../parquet/filter2/predicate/ParquetFilters.java  |  40 +++++-
 .../format/parquet/ParquetTypeWideningTest.java    | 151 +++++++++++++++++++++
 2 files changed, 188 insertions(+), 3 deletions(-)

diff --git 
a/paimon-format/src/main/java/org/apache/parquet/filter2/predicate/ParquetFilters.java
 
b/paimon-format/src/main/java/org/apache/parquet/filter2/predicate/ParquetFilters.java
index f7f5edbeeb..9f6936f6c8 100644
--- 
a/paimon-format/src/main/java/org/apache/parquet/filter2/predicate/ParquetFilters.java
+++ 
b/paimon-format/src/main/java/org/apache/parquet/filter2/predicate/ParquetFilters.java
@@ -543,6 +543,8 @@ public class ParquetFilters {
             throw new UnsupportedOperationException();
         }
 
+        validateNarrowIntegerPushdown(fieldRef, fileType);
+
         for (PrimitiveType.PrimitiveTypeName candidate : acceptable) {
             if (fileType.getPrimitiveTypeName() == candidate) {
                 return new PushdownTarget(fileType.getName(), candidate);
@@ -551,6 +553,32 @@ public class ParquetFilters {
         throw new UnsupportedOperationException();
     }
 
+    /** Reject physical integer ranges that the declared narrow type cannot 
preserve on read. */
+    private static void validateNarrowIntegerPushdown(FieldRef fieldRef, 
PrimitiveType fileType) {
+        int maxBitWidth;
+        switch (fieldRef.type().getTypeRoot()) {
+            case TINYINT:
+                maxBitWidth = 8;
+                break;
+            case SMALLINT:
+                maxBitWidth = 16;
+                break;
+            default:
+                return;
+        }
+
+        LogicalTypeAnnotation logicalType = 
fileType.getLogicalTypeAnnotation();
+        if (!(logicalType instanceof 
LogicalTypeAnnotation.IntLogicalTypeAnnotation)) {
+            throw new UnsupportedOperationException();
+        }
+
+        LogicalTypeAnnotation.IntLogicalTypeAnnotation intType =
+                (LogicalTypeAnnotation.IntLogicalTypeAnnotation) logicalType;
+        if (!intType.isSigned() || intType.getBitWidth() > maxBitWidth) {
+            throw new UnsupportedOperationException();
+        }
+    }
+
     /**
      * The file column a predicate has to name and the type it has to speak. 
The name matters in
      * case-insensitive mode: parquet-mr resolves a predicate against the file 
by exact column path,
@@ -571,8 +599,9 @@ public class ParquetFilters {
 
     /**
      * The physical types a predicate on this Paimon type can be expressed in, 
most preferred first.
-     * The head is what Paimon itself writes; the tail is the type widening 
the vectorized reader
-     * accepts, so pushdown stays available for exactly the files that can be 
read.
+     * The head is what Paimon itself writes; the tail is a different type the 
vectorized reader
+     * accepts without changing predicate semantics. Lossy narrowing is 
excluded even when the file
+     * can be read, because filtering happens before the reader converts the 
value.
      */
     private static PrimitiveType.PrimitiveTypeName[] acceptableTypes(
             org.apache.paimon.types.DataType type) {
@@ -583,6 +612,10 @@ public class ParquetFilters {
                 };
             case TINYINT:
             case SMALLINT:
+                // INT64 is lossy; INT32 annotations are validated by 
pushdownTarget.
+                return new PrimitiveType.PrimitiveTypeName[] {
+                    PrimitiveType.PrimitiveTypeName.INT32
+                };
             case INTEGER:
                 return new PrimitiveType.PrimitiveTypeName[] {
                     PrimitiveType.PrimitiveTypeName.INT32, 
PrimitiveType.PrimitiveTypeName.INT64
@@ -592,8 +625,9 @@ public class ParquetFilters {
                     PrimitiveType.PrimitiveTypeName.INT64, 
PrimitiveType.PrimitiveTypeName.INT32
                 };
             case FLOAT:
+                // A DOUBLE file value may round to the predicate's FLOAT 
value only after reading.
                 return new PrimitiveType.PrimitiveTypeName[] {
-                    PrimitiveType.PrimitiveTypeName.FLOAT, 
PrimitiveType.PrimitiveTypeName.DOUBLE
+                    PrimitiveType.PrimitiveTypeName.FLOAT
                 };
             case DOUBLE:
                 // A double bound cannot be narrowed to float without rounding 
it, so a FLOAT file
diff --git 
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetTypeWideningTest.java
 
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetTypeWideningTest.java
index d4d71f9014..04abccbf5a 100644
--- 
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetTypeWideningTest.java
+++ 
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetTypeWideningTest.java
@@ -221,6 +221,30 @@ class ParquetTypeWideningTest {
         assertThat(rows).hasSize(3);
     }
 
+    /** DOUBLE values are rounded while reading, so predicates cannot be 
pushed before the cast. */
+    @Test
+    void testDoubleReadAsFloatDoesNotPushLossyPredicate() throws Exception {
+        RowType readType =
+                RowType.builder()
+                        .field("pageviewId", DataTypes.STRING())
+                        .field("rate", DataTypes.FLOAT())
+                        .build();
+        Path path =
+                writeSchema(
+                        "message root {\n"
+                                + "  optional binary pageviewId (UTF8);\n"
+                                + "  optional double rate;\n"
+                                + "}",
+                        (group, i) ->
+                                group.append("pageviewId", 
PAGEVIEW_IDS[i]).append("rate", 0.1d));
+        PredicateBuilder builder = new PredicateBuilder(readType);
+
+        List<Object[]> rows =
+                read(readType, path, 
Collections.singletonList(builder.equal(1, 0.1f)));
+
+        assertThat(rows).extracting(row -> row[1]).containsExactly(0.1f, 0.1f, 
0.1f);
+    }
+
     // ------------------------------------------------------------------
     // Narrowing INT64 -> INT. The reader already handles it via
     // IntegerFromLongUpdater; only the pushdown was left behind.
@@ -243,6 +267,124 @@ class ParquetTypeWideningTest {
         assertThat((Integer) rows.get(2)[1]).isEqualTo(3000);
     }
 
+    /** The INT64 reader wraps values outside the TINYINT range, so pushdown 
would lose rows. */
+    @Test
+    void testInt64ReadAsTinyIntDoesNotPushLossyPredicate() throws Exception {
+        RowType readType =
+                RowType.builder()
+                        .field("pageviewId", DataTypes.STRING())
+                        .field("ecpm", DataTypes.TINYINT())
+                        .build();
+        Path path = write(ecpmSchema("int64 ecpm"), (group, i) -> 
group.append("ecpm", 128L));
+        PredicateBuilder builder = new PredicateBuilder(readType);
+
+        List<Object[]> rows =
+                read(readType, path, 
Collections.singletonList(builder.equal(1, (byte) -128)));
+
+        assertThat(rows)
+                .extracting(row -> row[1])
+                .containsExactly((byte) -128, (byte) -128, (byte) -128);
+    }
+
+    /** The INT64 reader wraps values outside the SMALLINT range, so pushdown 
would lose rows. */
+    @Test
+    void testInt64ReadAsSmallIntDoesNotPushLossyPredicate() throws Exception {
+        RowType readType =
+                RowType.builder()
+                        .field("pageviewId", DataTypes.STRING())
+                        .field("ecpm", DataTypes.SMALLINT())
+                        .build();
+        Path path = write(ecpmSchema("int64 ecpm"), (group, i) -> 
group.append("ecpm", 65_535L));
+        PredicateBuilder builder = new PredicateBuilder(readType);
+
+        List<Object[]> rows =
+                read(readType, path, 
Collections.singletonList(builder.equal(1, (short) -1)));
+
+        assertThat(rows)
+                .extracting(row -> row[1])
+                .containsExactly((short) -1, (short) -1, (short) -1);
+    }
+
+    /** A bare INT32 may also wrap when read as TINYINT, so its predicate 
cannot be pushed. */
+    @Test
+    void testInt32ReadAsTinyIntDoesNotPushLossyPredicate() throws Exception {
+        RowType readType =
+                RowType.builder()
+                        .field("pageviewId", DataTypes.STRING())
+                        .field("ecpm", DataTypes.TINYINT())
+                        .build();
+        Path path = write(ecpmSchema("int32 ecpm"), (group, i) -> 
group.append("ecpm", 128));
+        PredicateBuilder builder = new PredicateBuilder(readType);
+
+        List<Object[]> rows =
+                read(readType, path, 
Collections.singletonList(builder.equal(1, (byte) -128)));
+
+        assertThat(rows)
+                .extracting(row -> row[1])
+                .containsExactly((byte) -128, (byte) -128, (byte) -128);
+    }
+
+    /** A bare INT32 may also wrap when read as SMALLINT, so its predicate 
cannot be pushed. */
+    @Test
+    void testInt32ReadAsSmallIntDoesNotPushLossyPredicate() throws Exception {
+        RowType readType =
+                RowType.builder()
+                        .field("pageviewId", DataTypes.STRING())
+                        .field("ecpm", DataTypes.SMALLINT())
+                        .build();
+        Path path = write(ecpmSchema("int32 ecpm"), (group, i) -> 
group.append("ecpm", 65_535));
+        PredicateBuilder builder = new PredicateBuilder(readType);
+
+        List<Object[]> rows =
+                read(readType, path, 
Collections.singletonList(builder.equal(1, (short) -1)));
+
+        assertThat(rows)
+                .extracting(row -> row[1])
+                .containsExactly((short) -1, (short) -1, (short) -1);
+    }
+
+    /** A wider signed annotation still permits values outside the declared 
TINYINT range. */
+    @Test
+    void testInt16ReadAsTinyIntDoesNotPushLossyPredicate() throws Exception {
+        RowType readType =
+                RowType.builder()
+                        .field("pageviewId", DataTypes.STRING())
+                        .field("ecpm", DataTypes.TINYINT())
+                        .build();
+        Path path =
+                write(
+                        ecpmSchema("int32 ecpm (INTEGER(16,true))"),
+                        (group, i) -> group.append("ecpm", 128));
+        PredicateBuilder builder = new PredicateBuilder(readType);
+
+        List<Object[]> rows =
+                read(readType, path, 
Collections.singletonList(builder.equal(1, (byte) -128)));
+
+        assertThat(rows)
+                .extracting(row -> row[1])
+                .containsExactly((byte) -128, (byte) -128, (byte) -128);
+    }
+
+    /** Matching signed integer annotations retain predicate pushdown for 
valid physical values. */
+    @Test
+    void testAnnotatedInt8PushdownStillPrunes() throws Exception {
+        RowType readType =
+                RowType.builder()
+                        .field("pageviewId", DataTypes.STRING())
+                        .field("ecpm", DataTypes.TINYINT())
+                        .build();
+        Path path =
+                write(
+                        ecpmSchema("int32 ecpm (INTEGER(8,true))"),
+                        (group, i) -> group.append("ecpm", i + 1));
+        PredicateBuilder builder = new PredicateBuilder(readType);
+
+        List<Object[]> rows =
+                read(readType, path, 
Collections.singletonList(builder.greaterThan(1, (byte) 99)));
+
+        assertThat(rows).isEmpty();
+    }
+
     // ------------------------------------------------------------------
     // Control and mixed-file cases.
     // ------------------------------------------------------------------
@@ -406,12 +548,21 @@ class ParquetTypeWideningTest {
                 case VARCHAR:
                     values[i] = row.getString(i).toString();
                     break;
+                case TINYINT:
+                    values[i] = row.getByte(i);
+                    break;
+                case SMALLINT:
+                    values[i] = row.getShort(i);
+                    break;
                 case INTEGER:
                     values[i] = row.getInt(i);
                     break;
                 case BIGINT:
                     values[i] = row.getLong(i);
                     break;
+                case FLOAT:
+                    values[i] = row.getFloat(i);
+                    break;
                 case DOUBLE:
                     values[i] = row.getDouble(i);
                     break;

Reply via email to