andygrove opened a new pull request, #6154:
URL: https://github.com/apache/datafusion-comet/pull/6154
## Which issue does this PR close?
Closes #5992.
## Rationale for this change
This fixes #5992, plus a wrong-results bug in the same converter that turned
up while I was working
on it. The second one is the more serious of the two.
`icebergExprToProto` read a residual predicate's column name with
`term.ref().name()`. On an
`UnboundTransform`, `ref()` is the transform's source column, so a residual
of `bucket(4, id) = 2`
reached iceberg-rust as `id = 2`. iceberg-rust filters rows with the pushed
predicate, and the exact
filter above the scan can remove rows but cannot bring them back. With
Iceberg's SQL extensions
installed, which is what lets Spark push a system-function comparison into
the scan at all, this is
what `main` does:
```sql
CREATE TABLE t (id INT, data STRING) USING iceberg;
INSERT INTO t SELECT CAST(id AS INT), CONCAT('d', id) FROM range(100);
SELECT id FROM t WHERE system.bucket(4, id) = 2 AND data > 'a'; -- Comet 0
rows, Spark 22
SELECT id FROM t WHERE system.bucket(4, id) = 2 OR id = 50; -- Comet 1
row, Spark 23
```
`truncate` and `days` go wrong the same way, and the other transforms take
the same
`UnboundTransform` path. `CometScanRule` already falls back when the
residual is a bare transform
predicate, which is why a lone `WHERE system.bucket(4, id) = 2` is fine. It
does not look under
`AND`, `OR` or `NOT`, though, and adding any second predicate produces
exactly that shape. The same
`term.ref().name()` read is in the 1.0.0 tag, so this may be worth a
backport.
The other half is #5992. The converter and its call site in
`serializePartitions` both caught any
exception and carried on without the residual, so "this serde does not model
that predicate" and
"reflection against Iceberg's expression API broke" looked the same.
I want to flag one thing about #5992's premise. The review on #5515 that led
to it says a dropped
residual returns extra rows when Iceberg reports the filter as fully pushed.
I checked that against
`SparkScanBuilder.pushPredicates` (`BaseSparkScanBuilder` in 1.11.0) and the
residual logic it
relies on in all four pinned Iceberg versions (1.5.2, 1.8.1, 1.10.0 and
1.11.0), and I don't think
it can happen. Iceberg leaves a predicate out of `postScanFilters` only when
`ExpressionUtil.selectsPartitions` says it selects whole partitions in every
spec, meaning the
strict and inclusive projections agree. For exactly those predicates,
`ResidualEvaluator` settles
every file to `alwaysTrue` or `alwaysFalse`, so they never reach a residual.
Everything a residual
can hold is re-applied above the scan, and dropping one costs pruning, not
rows. I've still made
failures fatal, as the issue asks, because once `CometScanRule` has
committed the scan, failing is
the only way a misreading of Iceberg's API gets noticed. But that half of
the change is about
visibility rather than a known wrong answer, and I don't think #5992 on its
own is critical.
## What changes are included in this PR?
**Converter.** `icebergExprToProto` now accepts only a bare `NamedReference`
term and returns `None`
for anything else. That covers all six transforms, and `UnboundExtract` on
Iceberg versions that
have it. The function no longer catches exceptions. It returns `None` only
where it has decided not
to model something, and its scaladoc now says what the pushed predicate must
satisfy: it may be
weaker than the residual, but never stronger. The old scaladoc called it
only a pruning hint, but
iceberg-rust also filters rows with it. Most of this diff is re-indentation
from dropping the outer
`try`, and `?w=1` shows the real change.
**Call site.** `serializePartitions` goes through a new `private[operator]
serializeResidual`. It
turns a conversion failure into a `RuntimeException` that names the
residual, the same way
delete-file extraction already fails.
**Tests.** There is a new `CometIcebergResidualPushdownSuite`, registered in
both PR workflows. It
needs a session with Iceberg's extensions installed, and
`CometIcebergNativeSuite` doesn't have one.
Without them Spark keeps a system function as a `StaticInvoke`, which it
cannot push, so no
transform ever reaches a residual there.
## How are these changes tested?
`CometIcebergResidualPushdownSuite` runs five queries over an unpartitioned
Iceberg table, putting
`bucket`, `truncate` and `days` under `AND`, `OR` and `NOT`. It checks each
one against Spark and
asserts that the scan stayed native, so a fallback cannot make it pass
vacuously. Against `main`,
all five return wrong rows. With this change, all five match Spark.
The new tests in `CometIcebergNativeScanSuite` drive `serializeResidual`
with stand-ins for a
throwing accessor and a missing one, and expect the query to fail. They also
check that the
deliberate declines still come back as `None`, and that a transform term is
declined both on its own
and nested. Putting the swallowing `catch` back in either
`serializeResidual` or the converter turns
both failure tests red.
`CometIcebergNativeSuite`, `CometIcebergSystemFunctionSuite`,
`CometFuzzIcebergSuite` and
`IcebergReflectionSuite` pass locally on the default profile (Spark 4.1,
Iceberg 1.11), and the tree
compiles against Spark 3.5 with Scala 2.12. I've added `run-iceberg-tests`,
since this changes the
Iceberg scan serde.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]