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;