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)