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 ab5c255080 [format] Reject incompatible annotated integers as BIGINT 
(#9010)
ab5c255080 is described below

commit ab5c255080b20c168bde79f327e4de2710449d2e
Author: Dapeng Sun(孙大鹏) <[email protected]>
AuthorDate: Tue Aug 4 19:37:01 2026 +0800

    [format] Reject incompatible annotated integers as BIGINT (#9010)
---
 .../format/parquet/ParquetSchemaConverter.java     |  22 ++++
 .../reader/ParquetVectorUpdaterFactory.java        |  15 ++-
 .../parquet/filter2/predicate/ParquetFilters.java  |  11 ++
 .../paimon/format/parquet/ParquetFiltersTest.java  |  70 +++++++++++
 .../format/parquet/ParquetTypeWideningTest.java    | 140 +++++++++++++++++++++
 5 files changed, 255 insertions(+), 3 deletions(-)

diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetSchemaConverter.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetSchemaConverter.java
index 12da716ef4..309ed5c5ff 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetSchemaConverter.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetSchemaConverter.java
@@ -70,6 +70,28 @@ public class ParquetSchemaConverter {
                 && !((LogicalTypeAnnotation.IntLogicalTypeAnnotation) 
logicalType).isSigned();
     }
 
+    /**
+     * Whether an INT32 or INT64 column's logical annotation can be 
represented as BIGINT. Unsigned
+     * INT32 fits in BIGINT, but unsigned INT64 and non-integer annotations do 
not. Other physical
+     * types are outside this annotation check.
+     */
+    public static boolean isBigIntLogicalTypeCompatible(PrimitiveType type) {
+        PrimitiveType.PrimitiveTypeName physicalType = 
type.getPrimitiveTypeName();
+        if (physicalType != INT32 && physicalType != INT64) {
+            return true;
+        }
+
+        LogicalTypeAnnotation logicalType = type.getLogicalTypeAnnotation();
+        if (logicalType == null) {
+            return true;
+        }
+        if (!(logicalType instanceof 
LogicalTypeAnnotation.IntLogicalTypeAnnotation)) {
+            return false;
+        }
+        return physicalType == INT32
+                || ((LogicalTypeAnnotation.IntLogicalTypeAnnotation) 
logicalType).isSigned();
+    }
+
     /** Convert paimon {@link RowType} to parquet {@link MessageType}. */
     public static MessageType convertToParquetMessageType(RowType rowType) {
         return new MessageType(PAIMON_SCHEMA, convertToParquetTypes(rowType));
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 885d880514..6239c9fbb5 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
@@ -73,6 +73,7 @@ import java.nio.ByteOrder;
 import java.util.concurrent.TimeUnit;
 import java.util.function.Function;
 
+import static 
org.apache.paimon.format.parquet.ParquetSchemaConverter.isBigIntLogicalTypeCompatible;
 import static 
org.apache.paimon.format.parquet.ParquetSchemaConverter.isUnsignedInt;
 import static org.apache.paimon.utils.Preconditions.checkArgument;
 
@@ -182,12 +183,20 @@ public class ParquetVectorUpdaterFactory {
         @Override
         public UpdaterFactory visit(BigIntType bigIntType) {
             return c -> {
-                if (c.getPrimitiveType().getPrimitiveTypeName()
-                        == PrimitiveType.PrimitiveTypeName.INT32) {
+                PrimitiveType parquetType = c.getPrimitiveType();
+                if (!isBigIntLogicalTypeCompatible(parquetType)) {
+                    throw new UnsupportedOperationException(
+                            "Cannot read "
+                                    + parquetType.getPrimitiveTypeName()
+                                    + " logical type "
+                                    + parquetType.getLogicalTypeAnnotation()
+                                    + " as BIGINT");
+                }
+                if (parquetType.getPrimitiveTypeName() == 
PrimitiveType.PrimitiveTypeName.INT32) {
                     // The file kept the narrower int, either because the 
column was widened in
                     // the metastore after the data was written, or because it 
is unsigned and
                     // BIGINT is the only Paimon type that can hold every 
value.
-                    return isUnsignedInt(c.getPrimitiveType())
+                    return isUnsignedInt(parquetType)
                             ? new LongFromUnsignedIntegerUpdater()
                             : new LongFromIntegerUpdater();
                 }
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 9f6936f6c8..8cab825ca7 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
@@ -32,6 +32,7 @@ import org.apache.paimon.types.BinaryType;
 import org.apache.paimon.types.BlobType;
 import org.apache.paimon.types.BooleanType;
 import org.apache.paimon.types.CharType;
+import org.apache.paimon.types.DataTypeRoot;
 import org.apache.paimon.types.DataTypeVisitor;
 import org.apache.paimon.types.DateType;
 import org.apache.paimon.types.DecimalType;
@@ -537,6 +538,8 @@ public class ParquetFilters {
             return new PushdownTarget(fieldRef.name(), acceptable[0]);
         }
 
+        validateBigIntCompatibility(fieldRef, fileType);
+
         if (ParquetSchemaConverter.isUnsignedInt(fileType)) {
             // An unsigned column orders its statistics unsigned, so a signed 
bound would prune the
             // wrong row groups. The read still widens the column; only the 
pruning is given up.
@@ -553,6 +556,14 @@ public class ParquetFilters {
         throw new UnsupportedOperationException();
     }
 
+    /** Reject physical integer annotations that cannot be represented as 
BIGINT. */
+    private static void validateBigIntCompatibility(FieldRef fieldRef, 
PrimitiveType fileType) {
+        if (fieldRef.type().getTypeRoot() == DataTypeRoot.BIGINT
+                && 
!ParquetSchemaConverter.isBigIntLogicalTypeCompatible(fileType)) {
+            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;
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 1aff4e39e6..386e3d454a 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
@@ -99,6 +99,76 @@ class ParquetFiltersTest {
                 true);
     }
 
+    @Test
+    public void testBigIntPredicateIsNotPushedToDecimalInt32() {
+        RowType rowType =
+                new RowType(Collections.singletonList(new DataField(0, 
"long1", new BigIntType())));
+        MessageType schema =
+                new MessageType(
+                        "paimon_schema",
+                        Collections.singletonList(
+                                Types.optional(PrimitiveTypeName.INT32)
+                                        
.as(LogicalTypeAnnotation.decimalType(2, 9))
+                                        .named("long1")));
+        PredicateBuilder builder = new PredicateBuilder(rowType);
+
+        test(schema, builder.equal(0, 12345L), "", false);
+    }
+
+    @Test
+    public void testBigIntPredicateIsNotPushedToIncompatibleInt64() {
+        RowType rowType =
+                new RowType(Collections.singletonList(new DataField(0, 
"long1", new BigIntType())));
+        Predicate predicate = new PredicateBuilder(rowType).equal(0, 12345L);
+        List<PrimitiveType> incompatibleTypes =
+                Arrays.asList(
+                        Types.optional(PrimitiveTypeName.INT64)
+                                .as(LogicalTypeAnnotation.decimalType(2, 18))
+                                .named("long1"),
+                        Types.optional(PrimitiveTypeName.INT64)
+                                .as(
+                                        LogicalTypeAnnotation.timeType(
+                                                true, 
LogicalTypeAnnotation.TimeUnit.MICROS))
+                                .named("long1"),
+                        Types.optional(PrimitiveTypeName.INT64)
+                                .as(
+                                        LogicalTypeAnnotation.timestampType(
+                                                false, 
LogicalTypeAnnotation.TimeUnit.MICROS))
+                                .named("long1"),
+                        Types.optional(PrimitiveTypeName.INT64)
+                                .as(LogicalTypeAnnotation.intType(64, false))
+                                .named("long1"));
+
+        for (PrimitiveType type : incompatibleTypes) {
+            FilterCompat.Filter filter =
+                    ParquetFilters.convert(
+                            Collections.singletonList(predicate),
+                            new MessageType("paimon_schema", type),
+                            true);
+            assertThat(filter)
+                    .as("logical type %s", type.getLogicalTypeAnnotation())
+                    .isEqualTo(FilterCompat.NOOP);
+        }
+    }
+
+    /**
+     * INTEGER(64,true) means the same as an unannotated INT64, so the 
predicate is still pushed.
+     */
+    @Test
+    public void testBigIntPredicateIsPushedToSignedAnnotatedInt64() {
+        RowType rowType =
+                new RowType(Collections.singletonList(new DataField(0, 
"long1", new BigIntType())));
+        MessageType schema =
+                new MessageType(
+                        "paimon_schema",
+                        Types.optional(PrimitiveTypeName.INT64)
+                                .as(LogicalTypeAnnotation.intType(64, true))
+                                .named("long1"));
+        PredicateBuilder builder = new PredicateBuilder(rowType);
+
+        test(schema, builder.equal(0, 12345L), "eq(long1, 12345)", true);
+    }
+
     @Test
     public void testString() {
         RowType rowType =
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 04abccbf5a..214c836ec5 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
@@ -37,8 +37,11 @@ import 
org.apache.parquet.example.data.simple.SimpleGroupFactory;
 import org.apache.parquet.hadoop.ParquetWriter;
 import org.apache.parquet.hadoop.example.ExampleParquetWriter;
 import org.apache.parquet.hadoop.util.HadoopOutputFile;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
 import org.apache.parquet.schema.MessageType;
 import org.apache.parquet.schema.MessageTypeParser;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.apache.parquet.schema.Types;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
 
@@ -51,6 +54,7 @@ import java.util.UUID;
 import java.util.function.BiConsumer;
 
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
 /**
  * Reads where the declared type is wider than the type stored in the Parquet 
file.
@@ -169,6 +173,142 @@ class ParquetTypeWideningTest {
         assertThat(longs(rows, 1)).containsExactly(10L, 3_000_000_000L, 
4_294_967_295L);
     }
 
+    /** A DECIMAL annotation changes the meaning of the stored INT32 and 
cannot be widened. */
+    @Test
+    void testDecimalInt32CannotReadAsBigInt() throws Exception {
+        Path path =
+                write(
+                        ecpmSchema("int32 ecpm (DECIMAL(9,2))"),
+                        (group, i) -> group.append("ecpm", 12_345));
+
+        assertThatThrownBy(() -> read(ECPM_BIGINT, path, null))
+                .hasRootCauseInstanceOf(UnsupportedOperationException.class)
+                .rootCause()
+                .hasMessageContaining("DECIMAL");
+    }
+
+    /** A DATE annotation stores epoch days, not an integer that may be 
widened to BIGINT. */
+    @Test
+    void testDateInt32CannotReadAsBigInt() throws Exception {
+        Path path =
+                write(ecpmSchema("int32 ecpm (DATE)"), (group, i) -> 
group.append("ecpm", 20_000));
+
+        assertThatThrownBy(() -> read(ECPM_BIGINT, path, null))
+                .hasRootCauseInstanceOf(UnsupportedOperationException.class)
+                .rootCause()
+                .hasMessageContaining("DATE");
+    }
+
+    /** A TIME annotation stores milliseconds since midnight, not a generic 
integer. */
+    @Test
+    void testTimeInt32CannotReadAsBigInt() throws Exception {
+        MessageType schema =
+                new MessageType(
+                        "root",
+                        Arrays.asList(
+                                Types.optional(PrimitiveTypeName.BINARY)
+                                        .as(LogicalTypeAnnotation.stringType())
+                                        .named("pageviewId"),
+                                Types.optional(PrimitiveTypeName.INT32)
+                                        .as(
+                                                LogicalTypeAnnotation.timeType(
+                                                        true,
+                                                        
LogicalTypeAnnotation.TimeUnit.MILLIS))
+                                        .named("ecpm"),
+                                
Types.optional(PrimitiveTypeName.INT64).named("revenue")));
+        Path path =
+                writeGroups(
+                        schema,
+                        (group, i) ->
+                                group.append("pageviewId", PAGEVIEW_IDS[i])
+                                        .append("ecpm", 12_345)
+                                        .append("revenue", (long) i));
+
+        assertThatThrownBy(() -> read(ECPM_BIGINT, path, null))
+                .hasRootCauseInstanceOf(UnsupportedOperationException.class)
+                .rootCause()
+                .hasMessageContaining("TIME");
+    }
+
+    @Test
+    void testIncompatibleInt64CannotReadAsBigInt() throws Exception {
+        List<LogicalTypeAnnotation> incompatibleTypes =
+                Arrays.asList(
+                        LogicalTypeAnnotation.decimalType(2, 18),
+                        LogicalTypeAnnotation.timeType(true, 
LogicalTypeAnnotation.TimeUnit.MICROS),
+                        LogicalTypeAnnotation.timestampType(
+                                false, LogicalTypeAnnotation.TimeUnit.MICROS),
+                        LogicalTypeAnnotation.intType(64, false));
+
+        for (LogicalTypeAnnotation logicalType : incompatibleTypes) {
+            MessageType schema =
+                    new MessageType(
+                            "root",
+                            Arrays.asList(
+                                    Types.optional(PrimitiveTypeName.BINARY)
+                                            
.as(LogicalTypeAnnotation.stringType())
+                                            .named("pageviewId"),
+                                    Types.optional(PrimitiveTypeName.INT64)
+                                            .as(logicalType)
+                                            .named("ecpm"),
+                                    
Types.optional(PrimitiveTypeName.INT64).named("revenue")));
+            Path path =
+                    writeGroups(
+                            schema,
+                            (group, i) ->
+                                    group.append("pageviewId", PAGEVIEW_IDS[i])
+                                            .append("ecpm", 12_345L)
+                                            .append("revenue", (long) i));
+
+            assertThatThrownBy(() -> read(ECPM_BIGINT, path, null))
+                    .as("logical type %s", logicalType)
+                    
.hasRootCauseInstanceOf(UnsupportedOperationException.class)
+                    .rootCause()
+                    .hasMessageContaining("INT64")
+                    .hasMessageContaining("BIGINT");
+        }
+    }
+
+    /** INTEGER(64,true) means the same as an unannotated INT64 and must stay 
readable. */
+    @Test
+    void testSignedInt64AnnotationReadAsBigInt() throws Exception {
+        MessageType schema =
+                new MessageType(
+                        "root",
+                        Arrays.asList(
+                                Types.optional(PrimitiveTypeName.BINARY)
+                                        .as(LogicalTypeAnnotation.stringType())
+                                        .named("pageviewId"),
+                                Types.optional(PrimitiveTypeName.INT64)
+                                        .as(LogicalTypeAnnotation.intType(64, 
true))
+                                        .named("ecpm"),
+                                
Types.optional(PrimitiveTypeName.INT64).named("revenue")));
+        Path path =
+                writeGroups(
+                        schema,
+                        (group, i) ->
+                                group.append("pageviewId", PAGEVIEW_IDS[i])
+                                        .append("ecpm", ECPM_VALUES[i])
+                                        .append("revenue", (long) i));
+        PredicateBuilder builder = new PredicateBuilder(ECPM_BIGINT);
+
+        assertThat(longs(read(ECPM_BIGINT, path, null), 
1)).containsExactly(10L, 150L, 3000L);
+        assertThat(
+                        longs(
+                                read(
+                                        ECPM_BIGINT,
+                                        path,
+                                        
Collections.singletonList(builder.lessThan(1, 5000L))),
+                                1))
+                .containsExactly(10L, 150L, 3000L);
+        assertThat(
+                        read(
+                                ECPM_BIGINT,
+                                path,
+                                
Collections.singletonList(builder.greaterThan(1, 99999L))))
+                .isEmpty();
+    }
+
     // ------------------------------------------------------------------
     // FLOAT -> DOUBLE, the same hole in the other numeric family.
     // ------------------------------------------------------------------

Reply via email to