From cac76a8d3512962c709ef79c32b6f476581256af Mon Sep 17 00:00:00 2001
From: Xiangyi Zhu <82511136+zhuxiangyi@users.noreply.github.com>
Date: Sun, 30 Aug 2026 23:21:33 +0800
Subject: [PATCH] [core][format][spark] Support nested field predicate pushdown
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Predicates on a struct's sub-field are not pushed down today. `SparkExpressionConverter` rejects any `NamedReference` with more than one part, so `WHERE user.addr.city = 'Beijing'` is only evaluated by the engine after every row has been read.
This PR pushes such predicates down to the parquet row group / page level.
**Design.** Introduce `NestedFieldTransform`, a `Transform` holding the enclosing top-level `FieldRef` plus the ordered field names to descend into it. It is deliberately **not** a `FieldTransform`, so `LeafPredicate.fieldRefOptional()` stays empty for these predicates, and every consumer that equates a leaf with a top-level column — manifest stats evaluation, file index lookup, ORC pushdown, schema evolution rewriting, partition-only predicate detection — falls into its existing give-up path unchanged. That is why the diff contains no defensive guards at those call sites.
The path is stored as component names rather than positions and is re-resolved by name whenever the transform is remapped onto a different row type (as column pruning and row-level auth do): a leaf that was pruned away fails closed instead of silently resolving onto whatever now sits at that position, and a reordered row type still finds the original field.
The parquet side resolves the dotted name against the file schema and re-dispatches through the normal function visitor via the existing `visitNonFieldLeaf` hook, so every pushable function works on a nested field exactly as it does on a flat one, without per-function code. Decimal and timestamp predicates carry the file's own resolved path rather than the leaf's bare name, so they address the same column a flat predicate would.
**Supported.** `IS NULL`, `IS NOT NULL`, `=`, `<>`, `<`, `<=`, `>`, `>=`, `BETWEEN`, `IN`, `NOT IN`, and `AND`/`OR` mixing a nested field with a top-level one. Any nesting depth.
**Refused, falling back to engine evaluation:**
- any path component under a repeated group — parquet-mr cannot filter under repetition;
- a path descending into a non-row type;
- the enclosing column's name or any nested component containing a dot — parquet-mr addresses a column by a dot-joined path that cannot express such a name, so resolving it would either miss the real column or, on an unlucky schema, address a different one;
- everything the flat path already refuses (`startsWith` / `endsWith` / `contains` / `like`).
**Deliberately out of scope**, each independent of this change:
- *Manifest-level min/max and file index pruning.* `SimpleStats` is a positional row over top-level columns, so a nested leaf has no slot to read from — pushdown here is parquet-only. Making those layers work on nested fields requires reorganising statistics by field id, which is a separate and much larger change.
- *Schema evolution.* A data file whose schema version predates the table's current schema does not get this pushdown: `SchemaEvolutionUtil.devolveFilters` translates a predicate by its top-level `FieldRef`, which a nested predicate deliberately does not expose, so such predicates are dropped for those files. Results stay correct — the engine still evaluates the filter — but the parquet-level pruning is gone until a later write or compaction rewrites the files under the current schema.
- *ORC.* `OrcPredicateFunctionVisitor.visitNonFieldLeaf` returns empty and is untouched.
- *Flink.* `PredicateConverter` does not produce nested predicates today (a nested access arrives as a `GET` call, not a `FieldReferenceExpression`), so the Flink path never constructs a `NestedFieldTransform` and its behaviour is unchanged.
**Existing tables are unaffected.** No format change, no new option, read path only. Pruning uses row group and page statistics that are already present in existing files, so no rewrite or compaction is needed. Pushdown remains an optimisation: a non-partition data filter is also kept in `postScan`, so Spark still evaluates it row by row and the result set cannot change.
**Measured** on 400k rows with a wide struct, counting real bytes read:
| data layout | point `=` | `BETWEEN`, 1% | `BETWEEN`, 10% | absent value |
| --- | --- | --- | --- | --- |
| clustered on the nested field | 5.1% | 5.1% | 15.1% | 0.03% (footer only) |
| zone-ordered | 5.1% | 10.1% | 15.1% | 0.03% |
| randomly distributed | 100% | 100% | 100% | 0.03% |
A flat control column matched the nested column in every case. As with any min/max based pruning, the gain depends entirely on data locality.
This also adds the missing `PredicateBuilder.notIn(Transform, List)` overload — `notIn` was the only builder method without a `Transform` variant.
Tests reproduce every guarded scenario end to end against real row groups
(`ParquetFormatReadWriteTest`) as well as at the filter-construction and
transform level: nested predicates surviving remapping onto a pruned or
reordered row type, decimal predicates across every physical type parquet
can hold them in, both timestamp variants, and a nested or top-level
component whose name contains a dot. Run against Spark 3.3, 3.4 and 3.5.
No change to any on-disk format, no new table option. `NestedFieldTransform`
is registered as a `Transform` subtype so predicates serialise and
deserialise like the existing ones.
---
.../predicate/NestedFieldTransform.java | 196 +++++++++
.../paimon/predicate/PredicateBuilder.java | 4 +
.../apache/paimon/predicate/Transform.java | 1 +
.../predicate/NestedFieldTransformTest.java | 274 +++++++++++++
.../catalog/TableQueryAuthResultTest.java | 48 +++
.../filter2/predicate/ParquetFilters.java | 158 +++++--
.../format/parquet/ParquetFiltersTest.java | 386 ++++++++++++++++++
.../parquet/ParquetFormatReadWriteTest.java | 187 +++++++++
.../spark/util/SparkExpressionConverter.scala | 39 +-
.../sql/SparkV2FilterConverterTestBase.scala | 59 ++-
10 files changed, 1319 insertions(+), 33 deletions(-)
create mode 100644 paimon-common/src/main/java/org/apache/paimon/predicate/NestedFieldTransform.java
create mode 100644 paimon-common/src/test/java/org/apache/paimon/predicate/NestedFieldTransformTest.java
diff --git a/paimon-common/src/main/java/org/apache/paimon/predicate/NestedFieldTransform.java b/paimon-common/src/main/java/org/apache/paimon/predicate/NestedFieldTransform.java
new file mode 100644
index 000000000000..761de46119cf
--- /dev/null
+++ b/paimon-common/src/main/java/org/apache/paimon/predicate/NestedFieldTransform.java
@@ -0,0 +1,196 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.predicate;
+
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.RowType;
+
+import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonCreator;
+import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonIgnore;
+import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Objects;
+
+import static org.apache.paimon.utils.InternalRowUtils.get;
+import static org.apache.paimon.utils.Preconditions.checkArgument;
+
+/**
+ * Transform that extracts a field nested inside a row-typed column, for example {@code addr.city}.
+ *
+ *
The transform keeps the enclosing top-level column as its only {@link #inputs() input}, so
+ * anything that rewrites field indices (schema projection, for instance) keeps working without
+ * knowing about nesting. The positions below that column are held separately in {@link #path()}.
+ *
+ *
Deliberately not a {@link FieldTransform}: {@link LeafPredicate#fieldRefOptional()}
+ * returns empty for it, which is what keeps every consumer that equates a leaf with a top-level
+ * column — min/max pruning, file index lookup, ORC pushdown, schema evolution — from silently
+ * reading the enclosing column's metadata as if it belonged to the nested field. Those consumers
+ * give up on this transform instead, which costs pruning but never rows.
+ */
+public class NestedFieldTransform implements Transform {
+
+ private static final long serialVersionUID = 1L;
+
+ public static final String NAME = "NESTED_FIELD_REF";
+
+ public static final String FIELD_FIELD_REF = "fieldRef";
+ public static final String FIELD_PATH = "path";
+
+ /** The top-level row-typed column the nested field lives in. */
+ private final FieldRef fieldRef;
+
+ /**
+ * Names of the fields to descend into, relative to {@code fieldRef}'s row type. Never empty.
+ *
+ *
Deliberately names rather than positions: {@link #copyWithNewInputs} may be handed a
+ * structurally different row type — column masking and row filters remap that way — and a bare
+ * position would stay in range while silently addressing whatever now sits there. Names are
+ * re-resolved on every remap, so a reference either finds the same field or fails.
+ */
+ private final List path;
+
+ /** {@link #path} resolved to positions against {@code fieldRef}'s row type. */
+ private final int[] positions;
+
+ private final String name;
+ private final DataType outputType;
+
+ @JsonCreator
+ public NestedFieldTransform(
+ @JsonProperty(FIELD_FIELD_REF) FieldRef fieldRef,
+ @JsonProperty(FIELD_PATH) List path) {
+ checkArgument(path != null && !path.isEmpty(), "Nested field path must not be empty.");
+ this.fieldRef = fieldRef;
+ this.path = Collections.unmodifiableList(new ArrayList<>(path));
+ this.positions = new int[this.path.size()];
+
+ StringBuilder nameBuilder = new StringBuilder(fieldRef.name());
+ DataType current = fieldRef.type();
+ for (int i = 0; i < this.path.size(); i++) {
+ checkArgument(
+ current instanceof RowType,
+ "Nested field path of '%s' descends into a non-row type %s.",
+ fieldRef.name(),
+ current);
+ RowType rowType = (RowType) current;
+ String component = this.path.get(i);
+ int position = rowType.getFieldIndex(component);
+ checkArgument(
+ position >= 0,
+ "Nested field '%s' does not contain a field named '%s'.",
+ nameBuilder,
+ component);
+ positions[i] = position;
+ nameBuilder.append('.').append(component);
+ current = rowType.getTypeAt(position);
+ }
+ this.name = nameBuilder.toString();
+ this.outputType = current;
+ }
+
+ @Override
+ public String name() {
+ return NAME;
+ }
+
+ @JsonProperty(FIELD_FIELD_REF)
+ public FieldRef fieldRef() {
+ return fieldRef;
+ }
+
+ @JsonProperty(FIELD_PATH)
+ public List path() {
+ return path;
+ }
+
+ /** Dot-separated name from the top-level column down to the nested field, {@code addr.city}. */
+ @JsonIgnore
+ public String fieldName() {
+ return name;
+ }
+
+ @Override
+ @JsonIgnore
+ public List inputs() {
+ return Collections.singletonList(fieldRef);
+ }
+
+ @Override
+ @JsonIgnore
+ public DataType outputType() {
+ return outputType;
+ }
+
+ /**
+ * Reads the nested field out of {@code row}, which must match the row type {@link #fieldRef}
+ * was built against. A null anywhere along the path yields null, matching SQL semantics for
+ * field access on a null struct.
+ */
+ @Override
+ public Object transform(InternalRow row) {
+ int position = fieldRef.index();
+ if (row.isNullAt(position)) {
+ return null;
+ }
+ RowType currentType = (RowType) fieldRef.type();
+ InternalRow current = row.getRow(position, currentType.getFieldCount());
+
+ for (int i = 0; i < positions.length - 1; i++) {
+ position = positions[i];
+ if (current.isNullAt(position)) {
+ return null;
+ }
+ RowType nextType = (RowType) currentType.getTypeAt(position);
+ current = current.getRow(position, nextType.getFieldCount());
+ currentType = nextType;
+ }
+
+ int leaf = positions[positions.length - 1];
+ return get(current, leaf, currentType.getTypeAt(leaf));
+ }
+
+ @Override
+ public Transform copyWithNewInputs(List inputs) {
+ checkArgument(inputs.size() == 1);
+ return new NestedFieldTransform((FieldRef) inputs.get(0), path);
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (o == null || getClass() != o.getClass()) {
+ return false;
+ }
+ NestedFieldTransform that = (NestedFieldTransform) o;
+ return Objects.equals(fieldRef, that.fieldRef) && Objects.equals(path, that.path);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(fieldRef, path);
+ }
+
+ @Override
+ public String toString() {
+ return name;
+ }
+}
diff --git a/paimon-common/src/main/java/org/apache/paimon/predicate/PredicateBuilder.java b/paimon-common/src/main/java/org/apache/paimon/predicate/PredicateBuilder.java
index 04e813cdcff9..81221a3599bc 100644
--- a/paimon-common/src/main/java/org/apache/paimon/predicate/PredicateBuilder.java
+++ b/paimon-common/src/main/java/org/apache/paimon/predicate/PredicateBuilder.java
@@ -278,6 +278,10 @@ public Predicate notIn(int idx, List literals) {
return in(idx, literals).negate().get();
}
+ public Predicate notIn(Transform transform, List literals) {
+ return in(transform, literals).negate().get();
+ }
+
public Predicate between(int idx, Object includedLowerBound, Object includedUpperBound) {
DataField field = rowType.getFields().get(idx);
return new LeafPredicate(
diff --git a/paimon-common/src/main/java/org/apache/paimon/predicate/Transform.java b/paimon-common/src/main/java/org/apache/paimon/predicate/Transform.java
index ad01afcfb7be..687e489388f0 100644
--- a/paimon-common/src/main/java/org/apache/paimon/predicate/Transform.java
+++ b/paimon-common/src/main/java/org/apache/paimon/predicate/Transform.java
@@ -34,6 +34,7 @@
property = Transform.FIELD_NAME)
@JsonSubTypes({
@JsonSubTypes.Type(value = FieldTransform.class, name = FieldTransform.NAME),
+ @JsonSubTypes.Type(value = NestedFieldTransform.class, name = NestedFieldTransform.NAME),
@JsonSubTypes.Type(value = CastTransform.class, name = CastTransform.NAME),
@JsonSubTypes.Type(value = ConcatTransform.class, name = ConcatTransform.NAME),
@JsonSubTypes.Type(value = ConcatWsTransform.class, name = ConcatWsTransform.NAME),
diff --git a/paimon-common/src/test/java/org/apache/paimon/predicate/NestedFieldTransformTest.java b/paimon-common/src/test/java/org/apache/paimon/predicate/NestedFieldTransformTest.java
new file mode 100644
index 000000000000..91fe19297b03
--- /dev/null
+++ b/paimon-common/src/test/java/org/apache/paimon/predicate/NestedFieldTransformTest.java
@@ -0,0 +1,274 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.predicate;
+
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.GenericArray;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.JsonSerdeUtil;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Optional;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Test for {@link NestedFieldTransform}. */
+class NestedFieldTransformTest {
+
+ // user STRUCT>
+ private static final RowType ADDR_TYPE =
+ RowType.of(
+ new org.apache.paimon.types.DataType[] {DataTypes.STRING(), DataTypes.STRING()},
+ new String[] {"city", "zip"});
+ private static final RowType USER_TYPE =
+ RowType.of(
+ new org.apache.paimon.types.DataType[] {DataTypes.BIGINT(), ADDR_TYPE},
+ new String[] {"id", "addr"});
+ private static final RowType ROW_TYPE =
+ RowType.of(
+ new org.apache.paimon.types.DataType[] {DataTypes.INT(), USER_TYPE},
+ new String[] {"pk", "user"});
+
+ private static final FieldRef USER_REF = new FieldRef(1, "user", USER_TYPE);
+
+ private static GenericRow row(Object user) {
+ return GenericRow.of(1, user);
+ }
+
+ @Test
+ public void testReadOneLevel() {
+ NestedFieldTransform transform =
+ new NestedFieldTransform(USER_REF, Collections.singletonList("id"));
+
+ assertThat(transform.fieldName()).isEqualTo("user.id");
+ assertThat(transform.outputType()).isEqualTo(DataTypes.BIGINT());
+ assertThat(transform.transform(row(GenericRow.of(42L, null)))).isEqualTo(42L);
+ }
+
+ @Test
+ public void testReadTwoLevels() {
+ NestedFieldTransform transform =
+ new NestedFieldTransform(USER_REF, Arrays.asList("addr", "city"));
+
+ assertThat(transform.fieldName()).isEqualTo("user.addr.city");
+ assertThat(transform.outputType()).isEqualTo(DataTypes.STRING());
+
+ GenericRow addr =
+ GenericRow.of(
+ BinaryString.fromString("Beijing"), BinaryString.fromString("100080"));
+ assertThat(transform.transform(row(GenericRow.of(42L, addr))))
+ .isEqualTo(BinaryString.fromString("Beijing"));
+ }
+
+ /** The descent loop is recursive; three levels must read as well as two. */
+ @Test
+ public void testReadThreeLevels() {
+ RowType level3 =
+ RowType.of(
+ new org.apache.paimon.types.DataType[] {DataTypes.BIGINT()},
+ new String[] {"d"});
+ RowType level2 =
+ RowType.of(new org.apache.paimon.types.DataType[] {level3}, new String[] {"c"});
+ RowType level1 =
+ RowType.of(new org.apache.paimon.types.DataType[] {level2}, new String[] {"b"});
+ FieldRef ref = new FieldRef(0, "a", level1);
+
+ NestedFieldTransform transform =
+ new NestedFieldTransform(ref, Arrays.asList("b", "c", "d"));
+ assertThat(transform.fieldName()).isEqualTo("a.b.c.d");
+ assertThat(transform.outputType()).isEqualTo(DataTypes.BIGINT());
+
+ GenericRow row = GenericRow.of(GenericRow.of(GenericRow.of(GenericRow.of(42L))));
+ assertThat(transform.transform(row)).isEqualTo(42L);
+
+ // a null two levels down still yields null
+ GenericRow withNull = GenericRow.of(GenericRow.of(GenericRow.of((Object) null)));
+ assertThat(transform.transform(withNull)).isNull();
+ }
+
+ @Test
+ public void testNullAnywhereOnThePathYieldsNull() {
+ NestedFieldTransform transform =
+ new NestedFieldTransform(USER_REF, Arrays.asList("addr", "city"));
+
+ // the top-level column is null
+ assertThat(transform.transform(row(null))).isNull();
+ // an intermediate struct is null
+ assertThat(transform.transform(row(GenericRow.of(42L, null)))).isNull();
+ // the leaf itself is null
+ assertThat(transform.transform(row(GenericRow.of(42L, GenericRow.of(null, null)))))
+ .isNull();
+ }
+
+ @Test
+ public void testPredicateOnNullEvaluatesFalse() {
+ PredicateBuilder builder = new PredicateBuilder(ROW_TYPE);
+ Predicate predicate =
+ builder.equal(
+ new NestedFieldTransform(USER_REF, Arrays.asList("addr", "city")),
+ BinaryString.fromString("Beijing"));
+
+ assertThat(predicate.test(row(null))).isFalse();
+ assertThat(predicate.test(row(GenericRow.of(42L, null)))).isFalse();
+ }
+
+ /**
+ * The whole safety story rests on this: nothing that equates a leaf with a top-level column can
+ * mistake a nested field for one, because it never gets a {@link FieldRef} back.
+ */
+ @Test
+ public void testNoFieldRefIsExposed() {
+ LeafPredicate predicate =
+ (LeafPredicate)
+ new PredicateBuilder(ROW_TYPE)
+ .equal(
+ new NestedFieldTransform(
+ USER_REF, Collections.singletonList("id")),
+ 42L);
+
+ assertThat(predicate.fieldRefOptional()).isEmpty();
+ // the enclosing column is what schema-level rewrites see
+ assertThat(predicate.fieldNames()).containsExactly("user");
+ }
+
+ /** Min/max of the enclosing column say nothing about the nested field, so nothing is pruned. */
+ @Test
+ public void testStatsNeverPrune() {
+ Predicate predicate =
+ new PredicateBuilder(ROW_TYPE)
+ .equal(
+ new NestedFieldTransform(USER_REF, Collections.singletonList("id")),
+ 42L);
+
+ assertThat(
+ predicate.test(
+ 100L,
+ GenericRow.of(1, null),
+ GenericRow.of(10, null),
+ new GenericArray(new Object[] {0L, 0L})))
+ .isTrue();
+ }
+
+ @Test
+ public void testProjectionKeepsThePath() {
+ Predicate predicate =
+ new PredicateBuilder(ROW_TYPE)
+ .equal(
+ new NestedFieldTransform(USER_REF, Arrays.asList("addr", "city")),
+ 42L);
+
+ // "user" moves from index 1 to index 0
+ Optional projected =
+ predicate.visit(PredicateProjectionConverter.fromProjection(new int[] {1}));
+
+ assertThat(projected).isPresent();
+ NestedFieldTransform transform =
+ (NestedFieldTransform) ((LeafPredicate) projected.get()).transform();
+ assertThat(transform.fieldRef().index()).isEqualTo(0);
+ assertThat(transform.path()).containsExactly("addr", "city");
+ assertThat(transform.fieldName()).isEqualTo("user.addr.city");
+ }
+
+ @Test
+ public void testJsonRoundTrip() {
+ Predicate predicate =
+ new PredicateBuilder(ROW_TYPE)
+ .equal(
+ new NestedFieldTransform(USER_REF, Arrays.asList("addr", "city")),
+ BinaryString.fromString("Beijing"));
+
+ String json = JsonSerdeUtil.toJson(predicate);
+ assertThat(JsonSerdeUtil.fromJson(json, Predicate.class)).isEqualTo(predicate);
+ }
+
+ @Test
+ public void testRejectsPathThroughNonRowType() {
+ FieldRef arrayRef = new FieldRef(0, "tags", DataTypes.ARRAY(DataTypes.STRING()));
+ assertThatThrownBy(() -> new NestedFieldTransform(arrayRef, Collections.singletonList("x")))
+ .isInstanceOf(IllegalArgumentException.class);
+
+ assertThatThrownBy(() -> new NestedFieldTransform(USER_REF, Collections.emptyList()))
+ .isInstanceOf(IllegalArgumentException.class);
+ assertThatThrownBy(
+ () -> new NestedFieldTransform(USER_REF, Collections.singletonList("nope")))
+ .isInstanceOf(IllegalArgumentException.class);
+ }
+
+ /**
+ * Remapping must not let a nested reference drift onto a different field. Column pruning can
+ * hand {@code copyWithNewInputs} a structurally different row type — a bare position stays in
+ * range and silently addresses whatever now sits there. Row filters and column masks are
+ * remapped this way, so drifting has to fail closed rather than resolve elsewhere.
+ */
+ @Test
+ public void testRemapOntoAPrunedRowTypeDoesNotDrift() {
+ RowType full =
+ RowType.of(
+ new org.apache.paimon.types.DataType[] {
+ DataTypes.STRING(), DataTypes.STRING()
+ },
+ new String[] {"secret", "region"});
+ FieldRef infoRef = new FieldRef(0, "info", full);
+ NestedFieldTransform onSecret =
+ new NestedFieldTransform(infoRef, Collections.singletonList("secret"));
+ assertThat(onSecret.fieldName()).isEqualTo("info.secret");
+
+ // "secret" was pruned away; position 0 is now "region"
+ RowType pruned =
+ RowType.of(
+ new org.apache.paimon.types.DataType[] {DataTypes.STRING()},
+ new String[] {"region"});
+ FieldRef prunedRef = new FieldRef(0, "info", pruned);
+
+ assertThatThrownBy(() -> onSecret.copyWithNewInputs(Collections.singletonList(prunedRef)))
+ .isInstanceOf(IllegalArgumentException.class);
+ }
+
+ /** Remapping onto a reordered row type must keep addressing the same field. */
+ @Test
+ public void testRemapFollowsTheFieldWhenPositionsShift() {
+ RowType full =
+ RowType.of(
+ new org.apache.paimon.types.DataType[] {
+ DataTypes.STRING(), DataTypes.STRING()
+ },
+ new String[] {"secret", "region"});
+ NestedFieldTransform onSecret =
+ new NestedFieldTransform(
+ new FieldRef(0, "info", full), Collections.singletonList("secret"));
+
+ RowType reordered =
+ RowType.of(
+ new org.apache.paimon.types.DataType[] {
+ DataTypes.STRING(), DataTypes.STRING()
+ },
+ new String[] {"region", "secret"});
+ Transform remapped =
+ onSecret.copyWithNewInputs(
+ Collections.singletonList(new FieldRef(0, "info", reordered)));
+
+ assertThat(((NestedFieldTransform) remapped).fieldName()).isEqualTo("info.secret");
+ }
+}
diff --git a/paimon-core/src/test/java/org/apache/paimon/catalog/TableQueryAuthResultTest.java b/paimon-core/src/test/java/org/apache/paimon/catalog/TableQueryAuthResultTest.java
index 5b7f64fe8a59..71bd9af1ae22 100644
--- a/paimon-core/src/test/java/org/apache/paimon/catalog/TableQueryAuthResultTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/catalog/TableQueryAuthResultTest.java
@@ -24,7 +24,10 @@
import org.apache.paimon.predicate.FieldRef;
import org.apache.paimon.predicate.FieldTransform;
import org.apache.paimon.predicate.LeafPredicate;
+import org.apache.paimon.predicate.NestedFieldTransform;
import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.JsonSerdeUtil;
@@ -162,4 +165,49 @@ public void testHasRules() {
Map masking = Collections.singletonMap("display", maskJson());
assertThat(new TableQueryAuthResult(null, masking).hasRules()).isTrue();
}
+
+ private static RowType infoRowType(String... nestedFields) {
+ DataType[] types = new DataType[nestedFields.length];
+ for (int i = 0; i < types.length; i++) {
+ types[i] = DataTypes.STRING();
+ }
+ return RowType.of(
+ new DataType[] {DataTypes.INT(), RowType.of(types, nestedFields)},
+ new String[] {"pk", "info"});
+ }
+
+ private static Predicate rowFilterOnInfoSecret(RowType rowType) {
+ RowType info = (RowType) rowType.getTypeAt(1);
+ return new PredicateBuilder(rowType)
+ .equal(
+ new NestedFieldTransform(
+ new FieldRef(1, "info", info), Collections.singletonList("secret")),
+ org.apache.paimon.data.BinaryString.fromString("x"));
+ }
+
+ /**
+ * A row filter on a nested field must not silently follow column pruning onto a different
+ * field. Remapping resolves the components by name, so a pruned-away leaf fails closed rather
+ * than letting the policy address whatever now sits at that position.
+ */
+ @Test
+ void testNestedRowFilterDoesNotDriftWhenTheLeafIsPruned() {
+ Predicate filter = rowFilterOnInfoSecret(infoRowType("secret", "region"));
+
+ // the projection kept "info" but dropped "info.secret"
+ RowType pruned = infoRowType("region");
+ assertThatThrownBy(() -> TableQueryAuthResult.remapPredicate(filter, pruned))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("secret");
+ }
+
+ /** Remapping onto a reordered row type must keep addressing the same nested field. */
+ @Test
+ void testNestedRowFilterFollowsTheFieldWhenPositionsShift() {
+ Predicate filter = rowFilterOnInfoSecret(infoRowType("secret", "region"));
+
+ Predicate remapped =
+ TableQueryAuthResult.remapPredicate(filter, infoRowType("region", "secret"));
+ assertThat(remapped.toString()).contains("info.secret");
+ }
}
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 e28b8ee437e1..dee9d44df89c 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
@@ -25,6 +25,7 @@
import org.apache.paimon.predicate.FieldRef;
import org.apache.paimon.predicate.FunctionVisitor;
import org.apache.paimon.predicate.LeafPredicate;
+import org.apache.paimon.predicate.NestedFieldTransform;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.types.ArrayType;
import org.apache.paimon.types.BigIntType;
@@ -56,6 +57,7 @@
import org.apache.parquet.filter2.predicate.Operators.DoubleColumn;
import org.apache.parquet.filter2.predicate.Operators.FloatColumn;
import org.apache.parquet.io.api.Binary;
+import org.apache.parquet.schema.GroupType;
import org.apache.parquet.schema.LogicalTypeAnnotation;
import org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation;
import org.apache.parquet.schema.LogicalTypeAnnotation.TimestampLogicalTypeAnnotation;
@@ -77,6 +79,9 @@
/** Convert {@link Predicate} to {@link FilterCompat.Filter}. */
public class ParquetFilters {
+ /** Columns here are named, never indexed, so a nested field's index is left unset. */
+ private static final int UNUSED_INDEX = -1;
+
private ParquetFilters() {}
public static FilterCompat.Filter convert(
@@ -273,9 +278,31 @@ public FilterPredicate visitNotIn(FieldRef fieldRef, List literals) {
throw new UnsupportedOperationException();
}
+ /**
+ * A nested field carries no index into the file, only a path, so it is re-dispatched under
+ * a {@link FieldRef} naming that path. Every other transform - casts, string functions -
+ * has no column of its own to filter on and is given up here.
+ */
@Override
public FilterPredicate visitNonFieldLeaf(LeafPredicate predicate) {
- throw new UnsupportedOperationException();
+ if (!(predicate.transform() instanceof NestedFieldTransform)) {
+ throw new UnsupportedOperationException();
+ }
+ NestedFieldTransform nested = (NestedFieldTransform) predicate.transform();
+ // The path reaches parquet-mr as a dot-joined string, which it splits back into
+ // components. A component that itself contains a dot does not survive that round trip:
+ // the filter would address a column the file does not hold, and a missing column reads
+ // as all-null, pruning row groups that actually match. Give up the pruning instead.
+ if (nested.fieldRef().name().indexOf('.') >= 0) {
+ throw new UnsupportedOperationException();
+ }
+ for (String component : nested.path()) {
+ if (component.indexOf('.') >= 0) {
+ throw new UnsupportedOperationException();
+ }
+ }
+ FieldRef pathRef = new FieldRef(UNUSED_INDEX, nested.fieldName(), nested.outputType());
+ return predicate.function().visit(this, pathRef, predicate.literals());
}
private Set convertSets(List values, Class kclass, FieldRef fieldRef) {
@@ -306,7 +333,7 @@ private Comparable> toParquetObject(Object value, FieldRef fieldRef) {
DecimalType decimalType = (DecimalType) fieldRef.type();
Decimal decimal = normalizeDecimal((Decimal) value, decimalType);
PrimitiveType primitiveType =
- decimalPrimitiveType(fieldRef, fileSchema, caseSensitive);
+ decimalColumn(fieldRef, fileSchema, caseSensitive).type;
switch (primitiveType.getPrimitiveTypeName()) {
case INT32:
long intValue = toUnscaledLong(decimal);
@@ -327,7 +354,7 @@ private Comparable> toParquetObject(Object value, FieldRef fieldRef) {
if (value instanceof Timestamp) {
Timestamp timestamp = (Timestamp) value;
- timestampPrimitiveType(fieldRef, fileSchema, caseSensitive);
+ timestampColumn(fieldRef, fileSchema, caseSensitive);
int precision = getTimestampPrecision(type);
if (precision <= 3) {
// milliseconds
@@ -464,9 +491,10 @@ private Binary decimalToBinary(Decimal decimal, int numBytes) {
}
}
- private static PrimitiveType decimalPrimitiveType(
+ private static FileColumn decimalColumn(
FieldRef fieldRef, MessageType fileSchema, boolean caseSensitive) {
- PrimitiveType primitiveType = primitiveType(fieldRef, fileSchema, caseSensitive);
+ FileColumn column = fileColumn(fieldRef, fileSchema, caseSensitive);
+ PrimitiveType primitiveType = column.type;
LogicalTypeAnnotation logicalType = primitiveType.getLogicalTypeAnnotation();
if (!(logicalType instanceof DecimalLogicalTypeAnnotation)) {
throw new UnsupportedOperationException();
@@ -477,19 +505,20 @@ private static PrimitiveType decimalPrimitiveType(
if (decimalLogicalType.getScale() != ((DecimalType) fieldRef.type()).getScale()) {
throw new UnsupportedOperationException();
}
- return primitiveType;
+ return column;
}
- private static PrimitiveType timestampPrimitiveType(
+ private static FileColumn timestampColumn(
FieldRef fieldRef, MessageType fileSchema, boolean caseSensitive) {
- PrimitiveType primitiveType = primitiveType(fieldRef, fileSchema, caseSensitive);
+ FileColumn column = fileColumn(fieldRef, fileSchema, caseSensitive);
+ PrimitiveType primitiveType = column.type;
if (primitiveType.getPrimitiveTypeName() != PrimitiveType.PrimitiveTypeName.INT64) {
throw new UnsupportedOperationException();
}
LogicalTypeAnnotation logicalType = primitiveType.getLogicalTypeAnnotation();
if (logicalType == null) {
- return primitiveType;
+ return column;
}
if (!(logicalType instanceof TimestampLogicalTypeAnnotation)) {
throw new UnsupportedOperationException();
@@ -506,35 +535,100 @@ private static PrimitiveType timestampPrimitiveType(
|| timestampType.isAdjustedToUTC() != expectedAdjustedToUtc) {
throw new UnsupportedOperationException();
}
- return primitiveType;
+ return column;
}
- private static PrimitiveType primitiveType(
+ /**
+ * The column the file holds for {@code fieldRef}, with the file's own spelling of the path.
+ * Callers that build a parquet column must use {@link FileColumn#path}: a {@link PrimitiveType}
+ * only knows its own leaf name, so rebuilding the column from it drops the enclosing path and
+ * addresses a column the file does not have.
+ */
+ private static FileColumn fileColumn(
FieldRef fieldRef, MessageType fileSchema, boolean caseSensitive) {
- PrimitiveType matched = findPrimitiveType(fieldRef, fileSchema, caseSensitive);
+ FileColumn matched = findFileColumn(fieldRef, fileSchema, caseSensitive);
if (matched == null) {
throw new UnsupportedOperationException();
}
return matched;
}
+ /** A column the file actually holds: its own spelling of the path, and its physical type. */
+ private static class FileColumn {
+
+ private final String path;
+ private final PrimitiveType type;
+
+ private FileColumn(String path, PrimitiveType type) {
+ this.path = path;
+ this.type = type;
+ }
+ }
+
/**
* The file's column for {@code fieldRef}, or null when the file has no such column. A column
* that exists but is not primitive cannot carry a predicate at all, so it is rejected outright.
+ *
+ * {@code fieldRef} names a nested field with dots ({@code addr.city}), which is resolved by
+ * descending the file's groups. A top-level column matching the whole name wins over that walk,
+ * keeping flat columns spelled with dots resolving as they always did. parquet-mr identifies
+ * columns by dot-joined path too, so it cannot tell the two apart either way.
*/
@Nullable
- private static PrimitiveType findPrimitiveType(
+ private static FileColumn findFileColumn(
FieldRef fieldRef, MessageType fileSchema, boolean caseSensitive) {
- // Paimon predicates currently reference top-level fields only. Nested field
- // predicates are rejected before reaching the format reader.
- for (Type field : fileSchema.getFields()) {
+ Type matched = findChild(fileSchema, fieldRef.name(), caseSensitive);
+ if (matched != null) {
+ return toFileColumn(matched.getName(), matched);
+ }
+
+ String[] parts = fieldRef.name().split("\\.");
+ if (parts.length < 2) {
+ return null;
+ }
+
+ StringBuilder resolved = new StringBuilder();
+ GroupType parent = fileSchema;
+ for (int i = 0; i < parts.length; i++) {
+ Type child = findChild(parent, parts[i], caseSensitive);
+ if (child == null) {
+ return null;
+ }
+ if (child.getRepetition() == Type.Repetition.REPEATED) {
+ // A column under repetition has no one value per row, and parquet-mr refuses a
+ // predicate on it outright.
+ throw new UnsupportedOperationException();
+ }
+ if (i > 0) {
+ resolved.append('.');
+ }
+ resolved.append(child.getName());
+
+ if (i == parts.length - 1) {
+ return toFileColumn(resolved.toString(), child);
+ }
+ if (child.isPrimitive()) {
+ return null;
+ }
+ parent = child.asGroupType();
+ }
+ return null;
+ }
+
+ private static FileColumn toFileColumn(String path, Type field) {
+ if (!field.isPrimitive()) {
+ throw new UnsupportedOperationException();
+ }
+ return new FileColumn(path, field.asPrimitiveType());
+ }
+
+ @Nullable
+ private static Type findChild(GroupType parent, String name, boolean caseSensitive) {
+ for (Type field : parent.getFields()) {
if (caseSensitive
- ? field.getName().equals(fieldRef.name())
- : field.getName().equalsIgnoreCase(fieldRef.name())) {
- if (!field.isPrimitive()) {
- throw new UnsupportedOperationException();
- }
- return field.asPrimitiveType();
+ ? field.getName().equals(name)
+ : field.getName().equalsIgnoreCase(name)) {
+ return field;
}
}
return null;
@@ -557,10 +651,11 @@ private static PrimitiveType findPrimitiveType(
private static PushdownTarget pushdownTarget(
FieldRef fieldRef, MessageType fileSchema, boolean caseSensitive) {
PrimitiveType.PrimitiveTypeName[] acceptable = acceptableTypes(fieldRef.type());
- PrimitiveType fileType = findPrimitiveType(fieldRef, fileSchema, caseSensitive);
- if (fileType == null) {
+ FileColumn fileColumn = findFileColumn(fieldRef, fileSchema, caseSensitive);
+ if (fileColumn == null) {
return new PushdownTarget(fieldRef.name(), acceptable[0]);
}
+ PrimitiveType fileType = fileColumn.type;
validateBigIntCompatibility(fieldRef, fileType);
@@ -574,7 +669,7 @@ private static PushdownTarget pushdownTarget(
for (PrimitiveType.PrimitiveTypeName candidate : acceptable) {
if (fileType.getPrimitiveTypeName() == candidate) {
- return new PushdownTarget(fileType.getName(), candidate);
+ return new PushdownTarget(fileColumn.path, candidate);
}
}
throw new UnsupportedOperationException();
@@ -800,15 +895,16 @@ public Operators.Column> visit(TimeType timeType) {
@Override
public Operators.Column> visit(DecimalType decimalType) {
- PrimitiveType primitiveType = decimalPrimitiveType(fieldRef, fileSchema, caseSensitive);
+ FileColumn column = decimalColumn(fieldRef, fileSchema, caseSensitive);
+ PrimitiveType primitiveType = column.type;
switch (primitiveType.getPrimitiveTypeName()) {
case INT32:
- return FilterApi.intColumn(primitiveType.getName());
+ return FilterApi.intColumn(column.path);
case INT64:
- return FilterApi.longColumn(primitiveType.getName());
+ return FilterApi.longColumn(column.path);
case BINARY:
case FIXED_LEN_BYTE_ARRAY:
- return FilterApi.binaryColumn(primitiveType.getName());
+ return FilterApi.binaryColumn(column.path);
default:
throw new UnsupportedOperationException();
}
@@ -819,7 +915,7 @@ public Operators.Column> visit(TimestampType timestampType) {
int precision = timestampType.getPrecision();
if (precision <= 6) {
return FilterApi.longColumn(
- timestampPrimitiveType(fieldRef, fileSchema, caseSensitive).getName());
+ timestampColumn(fieldRef, fileSchema, caseSensitive).path);
}
// precision > 6 uses INT96, not supported for filter pushdown
throw new UnsupportedOperationException();
@@ -830,7 +926,7 @@ public Operators.Column> visit(LocalZonedTimestampType localZonedTimestampType
int precision = localZonedTimestampType.getPrecision();
if (precision <= 6) {
return FilterApi.longColumn(
- timestampPrimitiveType(fieldRef, fileSchema, caseSensitive).getName());
+ timestampColumn(fieldRef, fileSchema, caseSensitive).path);
}
// precision > 6 uses INT96, not supported for filter pushdown
throw new UnsupportedOperationException();
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 e6f4c4426046..65b936a185c6 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
@@ -18,13 +18,18 @@
package org.apache.paimon.format.parquet;
+import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.Decimal;
import org.apache.paimon.data.Timestamp;
+import org.apache.paimon.predicate.FieldRef;
+import org.apache.paimon.predicate.NestedFieldTransform;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.types.ArrayType;
import org.apache.paimon.types.BigIntType;
import org.apache.paimon.types.BooleanType;
import org.apache.paimon.types.DataField;
+import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DecimalType;
import org.apache.paimon.types.DoubleType;
import org.apache.paimon.types.FloatType;
@@ -1147,6 +1152,387 @@ private void test(
}
}
+ // ---------------------------------------------------------------------------------------
+ // nested fields
+ // ---------------------------------------------------------------------------------------
+
+ private static final RowType ADDR_TYPE =
+ RowType.of(
+ new DataType[] {new VarCharType(), new VarCharType()},
+ new String[] {"city", "zip"});
+
+ private static RowType nestedRowType() {
+ return RowType.of(
+ new DataType[] {
+ new BigIntType(),
+ RowType.of(
+ new DataType[] {new BigIntType(), ADDR_TYPE},
+ new String[] {"id", "addr"}),
+ new ArrayType(ADDR_TYPE)
+ },
+ new String[] {"pk", "user", "addrs"});
+ }
+
+ private static Predicate nestedPredicate(RowType rowType, String column, String... path) {
+ DataField field = rowType.getFields().get(rowType.getFieldIndex(column));
+ FieldRef ref = new FieldRef(rowType.getFieldIndex(column), field.name(), field.type());
+ return new PredicateBuilder(rowType)
+ .equal(
+ new NestedFieldTransform(ref, Arrays.asList(path)),
+ BinaryString.fromString("Beijing"));
+ }
+
+ @Test
+ public void testNestedField() {
+ RowType rowType = nestedRowType();
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+
+ // user.addr.city
+ test(
+ schema,
+ nestedPredicate(rowType, "user", "addr", "city"),
+ "eq(user.addr.city, Binary{\"Beijing\"})",
+ true);
+ }
+
+ /** A nested field is dispatched through the same visitors as a top-level one. */
+ @Test
+ public void testNestedFieldSupportsEveryPushableFunction() {
+ RowType rowType = nestedRowType();
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+ // user.id, a BIGINT one level down
+ FieldRef ref = new FieldRef(1, "user", rowType.getTypeAt(1));
+ NestedFieldTransform id = new NestedFieldTransform(ref, Collections.singletonList("id"));
+ PredicateBuilder builder = new PredicateBuilder(rowType);
+
+ test(schema, builder.isNull(id), "eq(user.id, null)", true);
+ test(schema, builder.isNotNull(id), "noteq(user.id, null)", true);
+ test(schema, builder.equal(id, 5L), "eq(user.id, 5)", true);
+ test(schema, builder.notEqual(id, 5L), "noteq(user.id, 5)", true);
+ test(schema, builder.lessThan(id, 5L), "lt(user.id, 5)", true);
+ test(schema, builder.lessOrEqual(id, 5L), "lteq(user.id, 5)", true);
+ test(schema, builder.greaterThan(id, 5L), "gt(user.id, 5)", true);
+ test(schema, builder.greaterOrEqual(id, 5L), "gteq(user.id, 5)", true);
+ test(schema, builder.between(id, 1L, 3L), "and(gteq(user.id, 1), lteq(user.id, 3))", true);
+ test(
+ schema,
+ builder.in(id, Arrays.asList(1L, 2L)),
+ "or(eq(user.id, 1), eq(user.id, 2))",
+ true);
+ test(
+ schema,
+ builder.notIn(id, Arrays.asList(1L, 2L)),
+ "and(noteq(user.id, 1), noteq(user.id, 2))",
+ true);
+
+ // AND/OR mixing a nested field with a top-level one
+ test(
+ schema,
+ PredicateBuilder.and(builder.greaterThan(id, 5L), builder.lessThan(0, 100L)),
+ "and(gt(user.id, 5), lt(pk, 100))",
+ true);
+
+ // string functions have no parquet equivalent, for nested and top-level alike
+ test(schema, builder.startsWith(id, BinaryString.fromString("x")), (String) null, false);
+ }
+
+ /**
+ * A field under a repeated group has no single value per row, and parquet-mr rejects a
+ * predicate on one outright. A table declaring a struct over a file that repeats it - a format
+ * table reading files someone else wrote - must give up rather than hand one over.
+ */
+ @Test
+ public void testNestedFieldUnderRepeatedGroupIsNotPushedDown() {
+ RowType rowType = nestedRowType();
+ MessageType schema =
+ new MessageType(
+ "paimon_schema",
+ Types.repeatedGroup()
+ .addField(Types.required(PrimitiveTypeName.INT64).named("id"))
+ .addField(
+ Types.requiredGroup()
+ .addField(
+ Types.required(PrimitiveTypeName.BINARY)
+ .as(
+ LogicalTypeAnnotation
+ .stringType())
+ .named("city"))
+ .named("addr"))
+ .named("user"));
+
+ test(schema, nestedPredicate(rowType, "user", "addr", "city"), (String) null, false);
+ }
+
+ /** A nested column the file does not hold still prunes: parquet-mr reads it as all-null. */
+ @Test
+ public void testNestedFieldMissingFromFile() {
+ RowType rowType = nestedRowType();
+ MessageType schema =
+ ParquetSchemaConverter.convertToParquetMessageType(
+ RowType.of(new DataType[] {new BigIntType()}, new String[] {"pk"}));
+
+ test(
+ schema,
+ nestedPredicate(rowType, "user", "addr", "city"),
+ "eq(user.addr.city, Binary{\"Beijing\"})",
+ true);
+ }
+
+ private static RowType payloadRowType() {
+ return RowType.of(
+ new DataType[] {
+ new BigIntType(),
+ RowType.of(
+ new DataType[] {
+ new DecimalType(10, 2),
+ new TimestampType(3),
+ new LocalZonedTimestampType(3),
+ new BigIntType()
+ },
+ new String[] {"amount", "ts", "ltz", "qty"}),
+ new DecimalType(10, 2)
+ },
+ new String[] {"pk", "payload", "amt_top"});
+ }
+
+ private static NestedFieldTransform payloadLeaf(RowType rowType, String leaf) {
+ RowType payload = (RowType) rowType.getTypeAt(1);
+ return new NestedFieldTransform(
+ new FieldRef(1, "payload", payload), Collections.singletonList(leaf));
+ }
+
+ /** Control: the two shapes that already worked must keep working. */
+ @Test
+ public void testTopLevelDecimalAndNestedBigIntAreUnaffected() {
+ RowType rowType = payloadRowType();
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+ PredicateBuilder builder = new PredicateBuilder(rowType);
+ Decimal amount = Decimal.fromBigDecimal(new BigDecimal("12.34"), 10, 2);
+
+ test(schema, builder.equal(2, amount), "eq(amt_top, 1234)", true);
+ test(
+ schema,
+ builder.greaterThan(payloadLeaf(rowType, "qty"), 5L),
+ "gt(payload.qty, 5)",
+ true);
+ }
+
+ /**
+ * A nested DECIMAL must be filtered on its full path. The physical type is resolved by walking
+ * the path, but the column handed to parquet-mr used to be rebuilt from the leaf {@code
+ * PrimitiveType}, which only knows its own name — parquet-mr then saw a missing top-level
+ * column and could drop every row group.
+ */
+ @Test
+ public void testNestedDecimalKeepsTheFullPath() {
+ RowType rowType = payloadRowType();
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+ PredicateBuilder builder = new PredicateBuilder(rowType);
+ NestedFieldTransform amount = payloadLeaf(rowType, "amount");
+ Decimal value = Decimal.fromBigDecimal(new BigDecimal("12.34"), 10, 2);
+
+ test(schema, builder.equal(amount, value), "eq(payload.amount, 1234)", true);
+ test(schema, builder.lessThan(amount, value), "lt(payload.amount, 1234)", true);
+ }
+
+ /** Same as {@link #testNestedDecimalKeepsTheFullPath()} for TIMESTAMP. */
+ @Test
+ public void testNestedTimestampKeepsTheFullPath() {
+ RowType rowType = payloadRowType();
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+ PredicateBuilder builder = new PredicateBuilder(rowType);
+ NestedFieldTransform ts = payloadLeaf(rowType, "ts");
+ Timestamp value = Timestamp.fromEpochMillis(1704067200000L);
+ long millis = value.getMillisecond();
+
+ test(schema, builder.equal(ts, value), "eq(payload.ts, " + millis + ")", true);
+ test(schema, builder.greaterThan(ts, value), "gt(payload.ts, " + millis + ")", true);
+ }
+
+ /** Same as {@link #testNestedDecimalKeepsTheFullPath()} for LOCAL ZONED TIMESTAMP. */
+ @Test
+ public void testNestedLocalZonedTimestampKeepsTheFullPath() {
+ RowType rowType = payloadRowType();
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+ PredicateBuilder builder = new PredicateBuilder(rowType);
+ NestedFieldTransform ltz = payloadLeaf(rowType, "ltz");
+ Timestamp value = Timestamp.fromEpochMillis(1704067200000L);
+ long millis = value.getMillisecond();
+
+ test(schema, builder.equal(ltz, value), "eq(payload.ltz, " + millis + ")", true);
+ }
+
+ /**
+ * A nested decimal schema with an explicit physical type, the way a Format Table may hold it.
+ */
+ private static MessageType nestedDecimalSchema(
+ PrimitiveTypeName physicalType, int fixedLength, int precision, int scale) {
+ Types.PrimitiveBuilder builder = Types.optional(physicalType);
+ if (physicalType == PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY) {
+ builder.length(fixedLength);
+ }
+ PrimitiveType amount =
+ builder.as(LogicalTypeAnnotation.decimalType(scale, precision)).named("amount");
+ return new MessageType(
+ "paimon_schema",
+ Arrays.asList(
+ Types.optional(PrimitiveTypeName.INT64).named("pk"),
+ Types.optionalGroup().addField(amount).named("payload")));
+ }
+
+ private void testNestedDecimalPhysicalType(
+ PrimitiveTypeName physicalType, int fixedLength, int precision, String literal) {
+ int scale = 2;
+ RowType payload =
+ RowType.of(
+ new DataType[] {new DecimalType(precision, scale)},
+ new String[] {"amount"});
+ RowType rowType =
+ RowType.of(
+ new DataType[] {new BigIntType(), payload}, new String[] {"pk", "payload"});
+ MessageType schema = nestedDecimalSchema(physicalType, fixedLength, precision, scale);
+ NestedFieldTransform amount =
+ new NestedFieldTransform(
+ new FieldRef(1, "payload", payload), Collections.singletonList("amount"));
+ Decimal value = Decimal.fromBigDecimal(new BigDecimal(literal), precision, scale);
+
+ FilterPredicate filter =
+ convert(schema, new PredicateBuilder(rowType).equal(amount, value));
+ // whatever the physical type, the column must be the full path
+ assertThat(filter.toString()).startsWith("eq(payload.amount, ");
+ }
+
+ /** The decimal visitor builds a column per physical type; every branch must keep the path. */
+ @Test
+ public void testNestedDecimalKeepsTheFullPathForEveryPhysicalType() {
+ // precision <= 9 -> INT32
+ testNestedDecimalPhysicalType(PrimitiveTypeName.INT32, 0, 8, "12.34");
+ // precision <= 18 -> INT64
+ testNestedDecimalPhysicalType(PrimitiveTypeName.INT64, 0, 15, "12.34");
+ // larger -> FIXED_LEN_BYTE_ARRAY
+ testNestedDecimalPhysicalType(PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY, 16, 30, "12.34");
+ // a Format Table may hold a decimal as BINARY
+ testNestedDecimalPhysicalType(PrimitiveTypeName.BINARY, 0, 30, "12.34");
+ }
+
+ /** Micros-precision timestamps take a different literal path than millis. */
+ @Test
+ public void testNestedTimestampMicrosKeepsTheFullPath() {
+ RowType payload =
+ RowType.of(
+ new DataType[] {new TimestampType(6), new LocalZonedTimestampType(6)},
+ new String[] {"ts", "ltz"});
+ RowType rowType =
+ RowType.of(
+ new DataType[] {new BigIntType(), payload}, new String[] {"pk", "payload"});
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+ PredicateBuilder builder = new PredicateBuilder(rowType);
+ FieldRef payloadRef = new FieldRef(1, "payload", payload);
+ Timestamp value = Timestamp.fromEpochMillis(1704067200000L);
+ long micros = value.toMicros();
+
+ test(
+ schema,
+ builder.equal(
+ new NestedFieldTransform(payloadRef, Collections.singletonList("ts")),
+ value),
+ "eq(payload.ts, " + micros + ")",
+ true);
+ test(
+ schema,
+ builder.greaterThan(
+ new NestedFieldTransform(payloadRef, Collections.singletonList("ltz")),
+ value),
+ "gt(payload.ltz, " + micros + ")",
+ true);
+ }
+
+ /** IN and NOT IN build the column through the same visitor; a nested decimal must keep it. */
+ @Test
+ public void testNestedDecimalInAndNotInKeepTheFullPath() {
+ RowType rowType = payloadRowType();
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+ PredicateBuilder builder = new PredicateBuilder(rowType);
+ NestedFieldTransform amount = payloadLeaf(rowType, "amount");
+ Decimal one = Decimal.fromBigDecimal(new BigDecimal("1.00"), 10, 2);
+ Decimal two = Decimal.fromBigDecimal(new BigDecimal("2.00"), 10, 2);
+
+ test(
+ schema,
+ builder.in(amount, Arrays.asList(one, two)),
+ "or(eq(payload.amount, 100), eq(payload.amount, 200))",
+ true);
+ test(
+ schema,
+ builder.notIn(amount, Arrays.asList(one, two)),
+ "and(noteq(payload.amount, 100), noteq(payload.amount, 200))",
+ true);
+ }
+
+ /** The path walk is recursive; three levels must resolve as well as two. */
+ @Test
+ public void testDeeplyNestedFieldKeepsTheFullPath() {
+ RowType level3 = RowType.of(new DataType[] {new BigIntType()}, new String[] {"d"});
+ RowType level2 = RowType.of(new DataType[] {level3}, new String[] {"c"});
+ RowType level1 = RowType.of(new DataType[] {level2}, new String[] {"b"});
+ RowType rowType =
+ RowType.of(new DataType[] {new BigIntType(), level1}, new String[] {"pk", "a"});
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+
+ NestedFieldTransform deep =
+ new NestedFieldTransform(
+ new FieldRef(1, "a", level1), Arrays.asList("b", "c", "d"));
+ test(schema, new PredicateBuilder(rowType).equal(deep, 7L), "eq(a.b.c.d, 7)", true);
+ }
+
+ /**
+ * A nested component whose own name contains a dot cannot be expressed as a dot-joined path:
+ * parquet-mr would split {@code s.a.b} into three components and miss the real two-component
+ * column, treating it as all-null and pruning matching row groups. Refuse the pushdown.
+ */
+ @Test
+ public void testNestedComponentContainingADotIsNotPushedDown() {
+ RowType inner = RowType.of(new DataType[] {new BigIntType()}, new String[] {"a.b"});
+ RowType rowType =
+ RowType.of(new DataType[] {new BigIntType(), inner}, new String[] {"pk", "s"});
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+ PredicateBuilder builder = new PredicateBuilder(rowType);
+
+ // the file really holds s -> "a.b"; a dot-joined "s.a.b" does not address it
+ assertThat(schema.getType("s").asGroupType().containsField("a.b")).isTrue();
+
+ NestedFieldTransform dotted =
+ new NestedFieldTransform(
+ new FieldRef(1, "s", inner), Collections.singletonList("a.b"));
+ test(schema, builder.equal(dotted, 7L), (String) null, false);
+ }
+
+ /**
+ * The dot may also sit in the top-level column's own name. The joined path then splits into
+ * components the file does not have — and, worse, could collide with a genuinely nested column
+ * of the same spelling. Refuse the pushdown here too.
+ */
+ @Test
+ public void testNestedFieldUnderATopLevelNameContainingADotIsNotPushedDown() {
+ RowType inner = RowType.of(new DataType[] {new VarCharType()}, new String[] {"city"});
+ RowType rowType =
+ RowType.of(new DataType[] {new BigIntType(), inner}, new String[] {"pk", "a.b"});
+ MessageType schema = ParquetSchemaConverter.convertToParquetMessageType(rowType);
+
+ // the file holds ["a.b", "city"]; the joined name "a.b.city" splits into [a, b, city]
+ assertThat(schema.containsField("a.b")).isTrue();
+ assertThat(schema.containsField("a")).isFalse();
+
+ NestedFieldTransform nested =
+ new NestedFieldTransform(
+ new FieldRef(1, "a.b", inner), Collections.singletonList("city"));
+ test(
+ schema,
+ new PredicateBuilder(rowType).equal(nested, BinaryString.fromString("Beijing")),
+ (String) null,
+ false);
+ }
+
private FilterPredicate convert(MessageType schema, Predicate predicate) {
FilterCompat.Filter filter =
ParquetFilters.convert(PredicateBuilder.splitAnd(predicate), schema, true);
diff --git a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFormatReadWriteTest.java b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFormatReadWriteTest.java
index 2784f157a113..d00dd68d9c84 100644
--- a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFormatReadWriteTest.java
+++ b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFormatReadWriteTest.java
@@ -19,6 +19,7 @@
package org.apache.paimon.format.parquet;
import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.Decimal;
import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericMap;
import org.apache.paimon.data.GenericRow;
@@ -36,7 +37,12 @@
import org.apache.paimon.format.SupportsWriterMetadata;
import org.apache.paimon.fs.PositionOutputStream;
import org.apache.paimon.options.Options;
+import org.apache.paimon.predicate.FieldRef;
+import org.apache.paimon.predicate.NestedFieldTransform;
+import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.predicate.PredicateBuilder;
import org.apache.paimon.reader.RecordReader;
+import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
@@ -51,7 +57,11 @@
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
+import java.io.IOException;
+import java.math.BigDecimal;
import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
@@ -292,4 +302,181 @@ public void testColumnCompressionCodec() throws Exception {
.containsEntry("name", CompressionCodecName.UNCOMPRESSED);
}
}
+
+ // -----------------------------------------------------------------------------------------
+ // end-to-end: a nested predicate must not silently drop rows
+ // -----------------------------------------------------------------------------------------
+
+ private RowType nestedPayloadType() {
+ return RowType.of(
+ new DataType[] {
+ DataTypes.BIGINT(),
+ RowType.of(
+ new DataType[] {
+ DataTypes.DECIMAL(10, 2),
+ DataTypes.TIMESTAMP(3),
+ DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE(3),
+ DataTypes.BIGINT()
+ },
+ new String[] {"amount", "ts", "ltz", "qty"})
+ },
+ new String[] {"pk", "payload"});
+ }
+
+ /**
+ * Reads with {@code predicate} pushed down and returns the primary keys that survived. Parquet
+ * filtering is row-group granular, so a matching row may come back alongside non-matching ones
+ * — what must never happen is the matching row disappearing.
+ */
+ private List readPks(RowType rowType, Predicate predicate) throws IOException {
+ List filters = new ArrayList<>();
+ filters.add(predicate);
+ List pks = new ArrayList<>();
+ try (RecordReader reader =
+ fileFormat()
+ .createReaderFactory(rowType, rowType, filters)
+ .createReader(
+ new FormatReaderContext(
+ fileIO, file, fileIO.getFileSize(file), null, null))) {
+ RecordReader.RecordIterator batch;
+ while ((batch = reader.readBatch()) != null) {
+ InternalRow row;
+ while ((row = batch.next()) != null) {
+ pks.add(row.getLong(0));
+ }
+ batch.releaseBatch();
+ }
+ }
+ return pks;
+ }
+
+ private void writeTwoPayloadRows(RowType rowType) throws IOException {
+ Decimal match = Decimal.fromBigDecimal(new BigDecimal("12.34"), 10, 2);
+ Decimal other = Decimal.fromBigDecimal(new BigDecimal("99.99"), 10, 2);
+ org.apache.paimon.data.Timestamp early =
+ org.apache.paimon.data.Timestamp.fromEpochMillis(1704067200000L);
+ org.apache.paimon.data.Timestamp late =
+ org.apache.paimon.data.Timestamp.fromEpochMillis(1704067200000L + 60_000L);
+ write(
+ fileFormat().createWriterFactory(rowType),
+ file,
+ GenericRow.of(1L, GenericRow.of(match, early, early, 7L)),
+ GenericRow.of(2L, GenericRow.of(other, late, late, 8L)));
+ }
+
+ private NestedFieldTransform payloadLeaf(RowType rowType, String leaf) {
+ RowType payload = (RowType) rowType.getTypeAt(1);
+ return new NestedFieldTransform(
+ new FieldRef(1, "payload", payload), Collections.singletonList(leaf));
+ }
+
+ /** Control: a nested BIGINT predicate already carried its full path and kept its row. */
+ @Test
+ public void testNestedBigIntPredicateKeepsMatchingRows() throws IOException {
+ RowType rowType = nestedPayloadType();
+ writeTwoPayloadRows(rowType);
+ Predicate onQty = new PredicateBuilder(rowType).equal(payloadLeaf(rowType, "qty"), 7L);
+ Assertions.assertThat(readPks(rowType, onQty))
+ .as("control: the row whose payload.qty equals 7 must survive the filter")
+ .contains(1L);
+ }
+
+ /**
+ * Reading with a predicate on a nested DECIMAL must still return the matching row. The column
+ * handed to parquet-mr used to carry only the leaf name, so parquet-mr saw a missing top-level
+ * column, treated it as all-null and pruned the row group holding the match.
+ */
+ @Test
+ public void testNestedDecimalPredicateKeepsMatchingRows() throws IOException {
+ RowType rowType = nestedPayloadType();
+ writeTwoPayloadRows(rowType);
+ Decimal match = Decimal.fromBigDecimal(new BigDecimal("12.34"), 10, 2);
+ Predicate onAmount =
+ new PredicateBuilder(rowType).equal(payloadLeaf(rowType, "amount"), match);
+ Assertions.assertThat(readPks(rowType, onAmount))
+ .as("the row whose payload.amount equals the literal must survive the filter")
+ .contains(1L);
+ }
+
+ /** Same as {@link #testNestedDecimalPredicateKeepsMatchingRows()} for LOCAL ZONED TIMESTAMP. */
+ @Test
+ public void testNestedLocalZonedTimestampPredicateKeepsMatchingRows() throws IOException {
+ RowType rowType = nestedPayloadType();
+ writeTwoPayloadRows(rowType);
+ org.apache.paimon.data.Timestamp early =
+ org.apache.paimon.data.Timestamp.fromEpochMillis(1704067200000L);
+ Predicate onLtz = new PredicateBuilder(rowType).equal(payloadLeaf(rowType, "ltz"), early);
+ Assertions.assertThat(readPks(rowType, onLtz))
+ .as("the row whose payload.ltz equals the literal must survive the filter")
+ .contains(1L);
+ }
+
+ /** Same as {@link #testNestedDecimalPredicateKeepsMatchingRows()} for TIMESTAMP. */
+ @Test
+ public void testNestedTimestampPredicateKeepsMatchingRows() throws IOException {
+ RowType rowType = nestedPayloadType();
+ writeTwoPayloadRows(rowType);
+ org.apache.paimon.data.Timestamp early =
+ org.apache.paimon.data.Timestamp.fromEpochMillis(1704067200000L);
+ Predicate onTs = new PredicateBuilder(rowType).equal(payloadLeaf(rowType, "ts"), early);
+ Assertions.assertThat(readPks(rowType, onTs))
+ .as("the row whose payload.ts equals the literal must survive the filter")
+ .contains(1L);
+ }
+
+ /**
+ * A nested component whose own name contains a dot cannot be addressed by a dot-joined path.
+ * Resolution fails, the filter is built against a path the file does not hold, and every row
+ * group is pruned — so the matching row disappears.
+ */
+ @Test
+ public void testNestedComponentContainingADotKeepsMatchingRows() throws IOException {
+ RowType inner = RowType.of(new DataType[] {DataTypes.BIGINT()}, new String[] {"a.b"});
+ RowType rowType =
+ RowType.of(new DataType[] {DataTypes.BIGINT(), inner}, new String[] {"pk", "s"});
+
+ write(
+ fileFormat().createWriterFactory(rowType),
+ file,
+ GenericRow.of(1L, GenericRow.of(7L)),
+ GenericRow.of(2L, GenericRow.of(8L)));
+
+ Predicate onDotted =
+ new PredicateBuilder(rowType)
+ .equal(
+ new NestedFieldTransform(
+ new FieldRef(1, "s", inner),
+ Collections.singletonList("a.b")),
+ 7L);
+
+ Assertions.assertThat(readPks(rowType, onDotted))
+ .as("the row whose s.`a.b` equals 7 must survive the filter")
+ .contains(1L);
+ }
+
+ /** The dot may sit in the top-level column's own name; the matching row must still survive. */
+ @Test
+ public void testTopLevelNameContainingADotKeepsMatchingRows() throws IOException {
+ RowType inner = RowType.of(new DataType[] {DataTypes.BIGINT()}, new String[] {"city"});
+ RowType rowType =
+ RowType.of(new DataType[] {DataTypes.BIGINT(), inner}, new String[] {"pk", "a.b"});
+
+ write(
+ fileFormat().createWriterFactory(rowType),
+ file,
+ GenericRow.of(1L, GenericRow.of(7L)),
+ GenericRow.of(2L, GenericRow.of(8L)));
+
+ Predicate onNested =
+ new PredicateBuilder(rowType)
+ .equal(
+ new NestedFieldTransform(
+ new FieldRef(1, "a.b", inner),
+ Collections.singletonList("city")),
+ 7L);
+
+ Assertions.assertThat(readPks(rowType, onNested))
+ .as("the row whose `a.b`.city equals 7 must survive the filter")
+ .contains(1L);
+ }
}
diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/SparkExpressionConverter.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/SparkExpressionConverter.scala
index 00898c98ca53..9d4fae2396d6 100644
--- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/SparkExpressionConverter.scala
+++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/SparkExpressionConverter.scala
@@ -61,7 +61,7 @@ object SparkExpressionConverter {
}
exp match {
- case n: NamedReference => Some(new FieldTransform(toPaimonFieldRef(n, rowType)))
+ case n: NamedReference => toPaimonFieldTransform(n, rowType)
case s: GeneralScalarExpression =>
s.name() match {
case CONCAT => convertChildren(s.children()).map(i => new ConcatTransform(i))
@@ -149,6 +149,43 @@ object SparkExpressionConverter {
}
}
+ /**
+ * A reference is either a top-level column or a path down into row-typed ones. Anything the path
+ * cannot descend - a field inside an array or a map, a name the schema does not hold - yields
+ * None, leaving the predicate for Spark to evaluate after the scan.
+ */
+ private def toPaimonFieldTransform(ref: NamedReference, rowType: RowType): Option[Transform] = {
+ val parts = ref.fieldNames()
+ val index = rowType.getFieldIndex(parts.head)
+ if (index == -1) {
+ return None
+ }
+ val root = rowType.getField(parts.head)
+ val rootRef = new FieldRef(index, root.name(), root.`type`())
+ if (parts.length == 1) {
+ return Some(new FieldTransform(rootRef))
+ }
+
+ // Keep the components Spark gave us: they are the transform's identity, and joining them
+ // would lose the boundaries of a name that itself contains a dot.
+ val path = new java.util.ArrayList[String](parts.length - 1)
+ var current = root.`type`()
+ parts.tail.foreach {
+ part =>
+ current match {
+ case nested: RowType =>
+ val position = nested.getFieldIndex(part)
+ if (position == -1) {
+ return None
+ }
+ path.add(part)
+ current = nested.getTypeAt(position)
+ case _ => return None
+ }
+ }
+ Some(new NestedFieldTransform(rootRef, path))
+ }
+
private def toPaimonFieldRef(ref: NamedReference, rowType: RowType): FieldRef = {
val fieldName = toFieldName(ref)
val f = rowType.getField(fieldName)
diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/SparkV2FilterConverterTestBase.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/SparkV2FilterConverterTestBase.scala
index ea80fe476d5d..7b2e53a06a3f 100644
--- a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/SparkV2FilterConverterTestBase.scala
+++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/SparkV2FilterConverterTestBase.scala
@@ -19,7 +19,7 @@
package org.apache.paimon.spark.sql
import org.apache.paimon.data.{BinaryString, Decimal, Timestamp}
-import org.apache.paimon.predicate.PredicateBuilder
+import org.apache.paimon.predicate.{FieldTransform, LeafPredicate, NestedFieldTransform, PredicateBuilder}
import org.apache.paimon.spark.{PaimonSparkTestBase, SparkV2FilterConverter}
import org.apache.paimon.spark.util.shim.TypeUtils.treatPaimonTimestampTypeAsSparkTimestampType
import org.apache.paimon.table.source.DataSplit
@@ -539,6 +539,63 @@ abstract class SparkV2FilterConverterTestBase extends PaimonSparkTestBase {
assert(filesScanned == 4, s"Expected 4 files but scanned $filesScanned files")
}
+ test("V2Filter: nested field") {
+ withTable("nested_tbl") {
+ sql("""
+ |CREATE TABLE nested_tbl (
+ | id INT,
+ | info STRUCT>
+ |) USING paimon
+ |""".stripMargin)
+ sql("INSERT INTO nested_tbl VALUES (1, struct(10, struct('Beijing', '100080')))")
+ sql("INSERT INTO nested_tbl VALUES (2, struct(20, struct('Shanghai', '200000')))")
+
+ val nestedConverter = SparkV2FilterConverter(loadTable("nested_tbl").rowType())
+
+ Seq("info.uid = 10" -> "info.uid", "info.addr.city = 'Beijing'" -> "info.addr.city")
+ .foreach {
+ case (filter, expectedName) =>
+ val predicate =
+ nestedConverter
+ .convert(v2Filter(filter, "nested_tbl"))
+ .get
+ .asInstanceOf[LeafPredicate]
+ val transform = predicate.transform().asInstanceOf[NestedFieldTransform]
+ assert(transform.fieldName() == expectedName)
+ // no FieldRef is handed out, so nothing mistakes this for a top-level column
+ assert(!predicate.fieldRefOptional().isPresent)
+ // the enclosing column is what field-name based rewrites see
+ assert(predicate.fieldNames().asScala == Seq("info"))
+
+ checkAnswer(sql(s"SELECT id FROM nested_tbl WHERE $filter"), Seq(Row(1)))
+ assert(
+ getPaimonScan(s"SELECT * FROM nested_tbl WHERE $filter").pushedDataFilters
+ .exists(_.toString.contains(expectedName)))
+ }
+
+ // a nested field still reads correctly alongside a projection of a sibling field
+ checkAnswer(
+ sql("SELECT info.addr.zip FROM nested_tbl WHERE info.addr.city = 'Shanghai'"),
+ Seq(Row("200000")))
+ }
+ }
+
+ test("V2Filter: a top-level column whose name contains a dot") {
+ withTable("dotted_tbl") {
+ sql("CREATE TABLE dotted_tbl (id INT, `a.b` STRING) USING paimon")
+
+ val dottedConverter = SparkV2FilterConverter(loadTable("dotted_tbl").rowType())
+ val predicate = dottedConverter
+ .convert(v2Filter("`a.b` = 'x'", "dotted_tbl"))
+ .get
+ .asInstanceOf[LeafPredicate]
+
+ // resolves as the flat column it is, not as a path into a struct named "a"
+ assert(predicate.transform().isInstanceOf[FieldTransform])
+ assert(predicate.fieldNames().asScala == Seq("a.b"))
+ }
+ }
+
private def v2Filter(str: String, tableName: String = "test_tbl"): SparkPredicate = {
val condition = sql(s"SELECT * FROM $tableName WHERE $str").queryExecution.optimizedPlan
.collectFirst { case f: Filter => f }