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]

Reply via email to