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 e4045f7100 [parquet] Avoid timestamp predicate pushdown for 
incompatible file schemas (#8797)
e4045f7100 is described below

commit e4045f7100a1ad5ed64362c9652b7a5c1c467953
Author: sanshi <[email protected]>
AuthorDate: Thu Jul 23 18:34:38 2026 +0800

    [parquet] Avoid timestamp predicate pushdown for incompatible file schemas 
(#8797)
---
 .../parquet/filter2/predicate/ParquetFilters.java  | 62 +++++++++++++++++-----
 .../paimon/format/parquet/ParquetFiltersTest.java  | 31 +++++++++++
 2 files changed, 81 insertions(+), 12 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 f6ffbed1e2..27d414ad75 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
@@ -56,6 +56,7 @@ import 
org.apache.parquet.filter2.predicate.Operators.FloatColumn;
 import org.apache.parquet.io.api.Binary;
 import org.apache.parquet.schema.LogicalTypeAnnotation;
 import 
org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation;
+import 
org.apache.parquet.schema.LogicalTypeAnnotation.TimestampLogicalTypeAnnotation;
 import org.apache.parquet.schema.MessageType;
 import org.apache.parquet.schema.PrimitiveType;
 import org.apache.parquet.schema.Type;
@@ -320,6 +321,7 @@ public class ParquetFilters {
                 return Binary.fromReusedByteArray((byte[]) value);
             } else if (value instanceof Timestamp) {
                 Timestamp timestamp = (Timestamp) value;
+                timestampPrimitiveType(fieldRef, fileSchema, caseSensitive);
                 int precision = getTimestampPrecision(type);
                 if (precision <= 3) {
                     // milliseconds
@@ -383,6 +385,51 @@ public class ParquetFilters {
 
     private static PrimitiveType decimalPrimitiveType(
             FieldRef fieldRef, MessageType fileSchema, boolean caseSensitive) {
+        PrimitiveType primitiveType = primitiveType(fieldRef, fileSchema, 
caseSensitive);
+        LogicalTypeAnnotation logicalType = 
primitiveType.getLogicalTypeAnnotation();
+        if (!(logicalType instanceof DecimalLogicalTypeAnnotation)) {
+            throw new UnsupportedOperationException();
+        }
+
+        DecimalLogicalTypeAnnotation decimalLogicalType =
+                (DecimalLogicalTypeAnnotation) logicalType;
+        if (decimalLogicalType.getScale() != ((DecimalType) 
fieldRef.type()).getScale()) {
+            throw new UnsupportedOperationException();
+        }
+        return primitiveType;
+    }
+
+    private static PrimitiveType timestampPrimitiveType(
+            FieldRef fieldRef, MessageType fileSchema, boolean caseSensitive) {
+        PrimitiveType primitiveType = primitiveType(fieldRef, fileSchema, 
caseSensitive);
+        if (primitiveType.getPrimitiveTypeName() != 
PrimitiveType.PrimitiveTypeName.INT64) {
+            throw new UnsupportedOperationException();
+        }
+
+        LogicalTypeAnnotation logicalType = 
primitiveType.getLogicalTypeAnnotation();
+        if (logicalType == null) {
+            return primitiveType;
+        }
+        if (!(logicalType instanceof TimestampLogicalTypeAnnotation)) {
+            throw new UnsupportedOperationException();
+        }
+
+        TimestampLogicalTypeAnnotation timestampType = 
(TimestampLogicalTypeAnnotation) logicalType;
+        int precision = getTimestampPrecision(fieldRef.type());
+        LogicalTypeAnnotation.TimeUnit expectedUnit =
+                precision <= 3
+                        ? LogicalTypeAnnotation.TimeUnit.MILLIS
+                        : LogicalTypeAnnotation.TimeUnit.MICROS;
+        boolean expectedAdjustedToUtc = fieldRef.type() instanceof 
LocalZonedTimestampType;
+        if (timestampType.getUnit() != expectedUnit
+                || timestampType.isAdjustedToUTC() != expectedAdjustedToUtc) {
+            throw new UnsupportedOperationException();
+        }
+        return primitiveType;
+    }
+
+    private static PrimitiveType primitiveType(
+            FieldRef fieldRef, MessageType fileSchema, boolean caseSensitive) {
         Type matched = null;
         // Paimon predicates currently reference top-level fields only. Nested 
field
         // predicates are rejected before reaching the format reader.
@@ -399,18 +446,7 @@ public class ParquetFilters {
             throw new UnsupportedOperationException();
         }
 
-        PrimitiveType primitiveType = matched.asPrimitiveType();
-        LogicalTypeAnnotation logicalType = 
primitiveType.getLogicalTypeAnnotation();
-        if (!(logicalType instanceof DecimalLogicalTypeAnnotation)) {
-            throw new UnsupportedOperationException();
-        }
-
-        DecimalLogicalTypeAnnotation decimalLogicalType =
-                (DecimalLogicalTypeAnnotation) logicalType;
-        if (decimalLogicalType.getScale() != ((DecimalType) 
fieldRef.type()).getScale()) {
-            throw new UnsupportedOperationException();
-        }
-        return primitiveType;
+        return matched.asPrimitiveType();
     }
 
     private static int getTimestampPrecision(org.apache.paimon.types.DataType 
type) {
@@ -523,6 +559,7 @@ public class ParquetFilters {
         public Operators.Column<?> visit(TimestampType timestampType) {
             int precision = timestampType.getPrecision();
             if (precision <= 6) {
+                timestampPrimitiveType(fieldRef, fileSchema, caseSensitive);
                 return FilterApi.longColumn(name);
             }
             // precision > 6 uses INT96, not supported for filter pushdown
@@ -533,6 +570,7 @@ public class ParquetFilters {
         public Operators.Column<?> visit(LocalZonedTimestampType 
localZonedTimestampType) {
             int precision = localZonedTimestampType.getPrecision();
             if (precision <= 6) {
+                timestampPrimitiveType(fieldRef, fileSchema, caseSensitive);
                 return FilterApi.longColumn(name);
             }
             // precision > 6 uses INT96, not supported for filter pushdown
diff --git 
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFiltersTest.java
 
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFiltersTest.java
index acb6e6e5bd..1aff4e39e6 100644
--- 
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFiltersTest.java
+++ 
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFiltersTest.java
@@ -882,6 +882,37 @@ class ParquetFiltersTest {
         test(schema, builder.greaterOrEqual(0, value), "gteq(ts1, " + 
expectedMicros + ")", true);
     }
 
+    @Test
+    public void testTimestampFileUnitMismatchCannotPushDown() {
+        Timestamp value = Timestamp.fromEpochMillis(1704067200123L, 456000);
+
+        RowType millisReadType =
+                new RowType(
+                        Collections.singletonList(new DataField(0, "ts1", new 
TimestampType(3))));
+        MessageType microsFileSchema =
+                new MessageType(
+                        "paimon_schema",
+                        Types.required(PrimitiveTypeName.INT64)
+                                .as(
+                                        LogicalTypeAnnotation.timestampType(
+                                                false, 
LogicalTypeAnnotation.TimeUnit.MICROS))
+                                .named("ts1"));
+        test(microsFileSchema, new PredicateBuilder(millisReadType).equal(0, 
value), "", false);
+
+        RowType microsReadType =
+                new RowType(
+                        Collections.singletonList(new DataField(0, "ts1", new 
TimestampType(6))));
+        MessageType millisFileSchema =
+                new MessageType(
+                        "paimon_schema",
+                        Types.required(PrimitiveTypeName.INT64)
+                                .as(
+                                        LogicalTypeAnnotation.timestampType(
+                                                false, 
LogicalTypeAnnotation.TimeUnit.MILLIS))
+                                .named("ts1"));
+        test(millisFileSchema, new PredicateBuilder(microsReadType).equal(0, 
value), "", false);
+    }
+
     @Test
     public void testLocalZonedTimestampMillis() {
         // precision <= 3 uses milliseconds (INT64)

Reply via email to