andygrove opened a new issue, #6562:
URL: https://github.com/apache/datafusion-comet/issues/6562
### Describe the bug
The native Iceberg writer counts NaNs that sit under a NULL parent in a data
file's `nan_value_counts`. Iceberg Java never sees those values, because a NULL
struct, list or map writes nothing for its children, so the native manifest can
record more NaNs for a field than the JVM writer does.
The NaN counter is iceberg-rust's `NanValueCountVisitor`
(`crates/iceberg/src/arrow/nan_val_cnt_visitor.rs` at Comet's pinned rev
`bb1e4a4`). It counts each float child array with only that child's own
validity. A struct child's value at a NULL struct slot is counted, and so are
the elements that a NULL list or map entry still points at. Arrow allows both
layouts. `nullif` produces them, because it replaces the top-level validity
buffer and leaves the children and offsets as they were.
A Comet plan produces them from `IF(cond, col, NULL)` or `CASE WHEN cond
THEN col END` over a nested column. For nested types Comet's `CaseWhenExpr`
evaluates through DataFusion's `CaseExpr`, which takes its
`InfallibleExprOrNull` path when the THEN branch is a column, and that path is
`nullif(col, NOT cond)`.
### Steps to reproduce
Write 100 rows where every third value is NaN to Parquet, then insert them
through a projection that nulls the odd rows:
```sql
CREATE TABLE t (id INT, v DOUBLE, s STRUCT<x: DOUBLE>, xs ARRAY<DOUBLE>)
USING iceberg;
-- src is a Parquet file with v = IF(id % 3 = 0, NaN, id), s =
named_struct('x', v), xs = array(v)
INSERT INTO t
SELECT id, IF(id % 2 = 0, v, NULL), IF(id % 2 = 0, s, NULL), IF(id % 2 = 0,
xs, NULL)
FROM src;
SELECT nan_value_counts FROM t.data_files;
```
The native plan is `CometIcebergWrite <- CometProject <- CometNativeScan`. I
compared it with the same insert through Iceberg Java:
| Field | 3.5 / Iceberg 1.8.1, native | JVM | 4.1 / Iceberg 1.11.0,
native | JVM |
| ------------ | --------------------------- | --- |
---------------------------- | ---- |
| `v` | 17 | 17 | 17
| 17 |
| `s.x` | 34 | 17 | 34
| 17 |
| `xs.element` | 34 | 17 | none
| none |
The top-level column is right, because `nullif` sets the double array's own
validity. The struct field is wrong on every profile: from Iceberg 1.10,
`ParquetMetrics` drops the metrics of fields under a list or map but keeps
those of struct fields, and Comet passes the native NaN count through. The list
element is wrong only where nested metrics are kept, Iceberg before 1.10 (the
Spark 3.4 and 3.5 profiles). Map values go through the same path as list
elements.
Both tables return the same rows, and the Parquet data is right, since the
Arrow writer honors the parent's validity. Only the metric is wrong. Spark can
push predicates onto struct fields, though, and a NaN count that reaches the
field's value count would make Iceberg's metrics evaluators treat the file as
holding nothing but NaNs. I haven't built a case where that changes a query
result.
### Expected behavior
`nan_value_counts` matches Iceberg Java: a value under a NULL struct, list
or map is not counted.
### Additional context
This is the same visitor as #6146, where it counts the list and map elements
outside an already-sliced batch. #6238 fixes that by gathering such a batch
with `take`, but that doesn't cover this:
- A batch whose children span its rows, as in the reproduction, is handed on
without a gather.
- arrow 59.3's `take` drops the elements of a NULL list or map entry at the
level it gathers. It keeps a NULL struct slot's child values, though, since it
takes every child at every index. It also keeps the elements under NULL entries
of a list nested inside the gathered one, which `MutableArrayData` copies as
whole ranges.
A Comet-side fix could push each struct's validity into its children and
drop the ranges of NULL list and map entries before the writer counts,
recursively. Alternatively the visitor could honor the parent's validity and
offsets upstream, which would also fix #6146 and let both workarounds go.
The same 4.1 run also shows a smaller, separate difference: `value_counts`
and `null_value_counts` for `s.x` are 100 and 50 natively but 50 and 0 through
the JVM. This looks like Iceberg Java 1.10+ taking a float field's metrics from
its value writers, which skip a field under a NULL struct, while the Parquet
footer counts it as a null. The native writer and Iceberg 1.8.1 both report the
footer's counts.
Part of #5649. Found while addressing review feedback on #6238.
--
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]