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]