zhuxiangyi opened a new pull request, #10144:
URL: https://github.com/apache/paimon/pull/10144

   ### Purpose
   
   Follow-up to #9423, which added predicate pushdown for fields nested inside 
a row and wired it up for Spark. Flink was left out: a predicate on a nested 
field is still evaluated entirely by Flink, and Paimon reads every file and 
every row for it.
   
   Flink does hand the predicate to the source. 
`PushFilterIntoTableSourceScanRule` goes through 
`FlinkRexUtil.extractPredicates`, which builds its 
`RexNodeToExpressionConverter` with the row type, so `visitFieldAccess` 
resolves nested access to a `NestedFieldReferenceExpression`. Only the child 
type differs from a top-level predicate:
   
   ```
   WHERE s.a = 7   ->  CallExpression equals(`s.a`, 7)
                         child: NestedFieldReferenceExpression `s.a`
                         child: ValueLiteralExpression 7
   
   WHERE t = 10    ->  CallExpression equals(t, 10)
                         child: FieldReferenceExpression t
                         child: ValueLiteralExpression 10
   ```
   
   `PredicateConverter` recognised only `FieldReferenceExpression`, so the 
nested one raised `UnsupportedExpression` and the whole filter went 
unconverted. That is what this PR fixes.
   
   `NestedFieldReferenceExpression.getFieldNames()` is the path as the names of 
the rows walked through, which is exactly how `NestedFieldTransform` addresses 
a nested field, so the two map onto each other directly. `ResolvedField` now 
carries either a top-level field index or a transform, and each operation picks 
the matching `PredicateBuilder` overload.
   
   Covered: the six comparisons (either operand order), `IN` / `NOT IN`, `IS 
(NOT) NULL`, `BETWEEN` / `NOT BETWEEN`, a `LIKE` prefix, `IS TRUE` / `IS FALSE` 
/ `IS NOT TRUE` / `IS NOT FALSE`, and combinations with top-level predicates. 
The existing guards extend to nested fields as well: negated floating-point 
comparisons stay unconverted, because Flink compares FLOAT/DOUBLE with Java 
operators while Paimon uses `compareTo`, and `NOT IN (..., NULL, ...)` is still 
always false.
   
   A path that does not address a field of the table - the root is not in the 
schema, the root is not a row, or a field along the way was renamed or dropped 
- stays unconverted and is left for Flink to evaluate. Results cannot change 
either way: `applyFilters` only consumes a filter that is a predicate over 
partition keys alone, which a nested one never is, so it always comes back in 
`unConsumedFilters` and Flink re-evaluates it after the scan. Pushdown only 
decides how much is read.
   
   **Flink version compatibility.** `NestedFieldReferenceExpression` only 
exists from Flink 1.19 on, but `paimon-flink-common` is compiled once and 
shaded into the bundles for Flink 1.16 to 1.18 as well. Every reference to it 
is therefore isolated in `NestedFieldReferences.Holder`, a class that is only 
loaded after `Class.forName` has confirmed the expression is on the classpath; 
without that isolation the `instanceof` resolves the missing class and throws 
`NoClassDefFoundError` at pushdown time.
   
   ### Tests
   
   - `NestedPredicateConverterTest` - the converted `Predicate` for every 
supported operation, the field on either side of a comparison, the negated 
forms, the floating-point and `NULL` guards, a path through several rows, and 
the paths that must stay unconverted. 12 of its 17 cases fail without the 
change; the other 5 are the negative ones, which hold either way.
   - `NestedFieldReferencesTest` - loads `NestedFieldReferences` through a 
class loader that hides `NestedFieldReferenceExpression`, standing in for Flink 
1.16 to 1.18, and asserts it answers rather than fails to link. Removing the 
guard makes it fail with the `NoClassDefFoundError` above.
   - `NestedFieldFilterPushDownITCase` - end to end over real data: matching 
rows, `NULL` rows, negated predicates and combinations with top-level 
predicates all return what they did before.
   
   The plan is deliberately not asserted on: `applyFilters` reports every 
filter as accepted whether or not it could be converted, so `filter=[...]` in 
the scan digest looks the same before and after and cannot tell the two apart.
   
   Run against Flink 1.20 and, with `-Pflink2`, Flink 2.2: 127 tests each, 
including `PredicateConverterTest`, `FilterPushDownITCase`, 
`FilterPushdownWithSchemaChangeITCase` and `SimpleSqlPredicateConvertorTest` 
for regressions.
   
   ### Related
   
   - #9423 - added `NestedFieldTransform`, the parquet-side pruning and the 
Spark converter. This PR is its Flink counterpart: all of the pruning machinery 
already exists, so outside `paimon-flink-common` the only change is one 
accessor on `PredicateBuilder`.
   
   With Spark covered by #9423 and Flink by this PR, the remaining gap is ORC: 
`OrcPredicateFunctionVisitor.visitNonFieldLeaf` declines every transform 
outright, so a nested predicate prunes on parquet but not on ORC. That is 
independent of this change and would be a separate PR.
   
   ### API and Format
   
   `PredicateBuilder.rowType()` is added, an additive accessor on an existing 
`@Public` class. No format change.
   
   ### Documentation
   
   No changes.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to