andygrove opened a new pull request, #6238:
URL: https://github.com/apache/datafusion-comet/pull/6238

   ## Which issue does this PR close?
   
   Closes #6146.
   
   ## Rationale for this change
   
   iceberg-rust's NaN counter reaches list elements through 
`list_array.values()` and map entries through `map_array.entries()`. Both 
return the whole child array and ignore the parent's offsets. The writer's 
`RowSlicer` already gathers the ranges it cuts itself, but it passes a whole 
input batch through unchanged. So when a batch arrives already sliced, its list 
and map children still hold the values of rows outside it, and their NaNs are 
counted too.
   
   The impact is narrower than the issue suggests:
   
   - It only reaches the manifest on Iceberg 1.5.2 and 1.8.1, the Spark 3.4 and 
3.5 profiles. From 1.10, iceberg-java builds metrics with `ParquetMetrics`, 
which keeps none for a field under a list or map. Comet rebuilds manifest 
metrics through that same code, so on the 4.x profiles neither writer records 
nested NaN counts.
   - A plain `LIMIT` doesn't reproduce it. The writer's input crosses Arrow 
Java (`ColumnarBatchArrowReader`), which cuts a list's child at the batch's 
last offset. That drops the rows after a slice but not the rows before one. A 
slice that starts past the first row does reach the writer, and `LIMIT ... 
OFFSET` produces one. On Spark 3.5, `coalesce(1).offset(40).limit(25)` recorded 
22 NaNs per nested field on the native path, and 8 through iceberg-java.
   - The result is wrong metadata, not wrong query results, since Spark cannot 
push predicates onto list elements or map values.
   
   ## What changes are included in this PR?
   
   - `RowSlicer::compact` runs on every input batch after the field-id cast. It 
gathers the batch whole with the slicer's existing `gather_rows` when two 
things hold: the schema has a float under a list or map (the slicer's existing 
`gather` decision), and some list or map in the batch has a child that reaches 
past the batch's rows. Otherwise it returns the batch unchanged, so a compact 
batch costs only the offset check.
   - Checking at entry covers all three writers. The unpartitioned pacer, the 
clustered splitter and iceberg-rust's fanout splitter can each pass a 
single-partition batch through whole. The fanout splitter does it through 
arrow's all-true filter shortcut.
   - `floats_outside_window` does the check. It compares each list's or map's 
first and last offsets with its child's length, and recurses through structs 
and nested lists. A container it doesn't inspect counts as reaching past 
whenever it holds a float.
   
   Not covered: the counter would also count NaNs stored under null list 
entries that have a non-empty range, and `take` keeps those ranges. I haven't 
found a plan that produces such lists.
   
   ## How are these changes tested?
   
   - Rust:
     - `nan_counts_ignore_rows_outside_an_already_sliced_batch` writes a 
pre-sliced batch with `list<double>` and `map<string, double>` columns through 
the unpartitioned, fanout and clustered writers. It uses two windows, a leading 
one and one that leaves rows out at both ends. On main, every case records 4 
NaNs where 1 or 2 are expected.
     - Unit tests cover `floats_outside_window` and `compact`. The detection 
cases are a list sliced from either end, a sliced struct holding a list, a 
nested list inside a compact outer list, and a list with no float. `compact` 
hands a compact batch through as the same arrays, and gathers a sliced one with 
its values unchanged.
     - Removing the call site, the last-offset check, or the recursion into a 
list's child each makes a test fail.
   - JVM: the new `CometIcebergWriteActionSuite` test "native acceleration: NaN 
counts under an OFFSET skip the list and map rows it drops" writes 
`coalesce(1).offset(40).limit(25)` through both writers and compares 
`nan_value_counts`.
     - It fails on main on Spark 3.5 (22 against 8), and passes with the fix on 
3.4 and 3.5.
     - From Iceberg 1.10 it can only check that the two writers agree, so it 
can't catch this bug on the default 4.1 profile. The Rust tests can.
     - It uses `coalesce(1)` because `INSERT ... SELECT ... LIMIT` over a file 
scan puts a shuffle before the limit, and the shuffle compacts the batch.
   - The full `CometIcebergWriteActionSuite` passes on Spark 4.1 (75 tests). 
`cargo clippy --all-targets --workspace -- -D warnings` and `cargo fmt` are 
clean.
   


-- 
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