andygrove commented on issue #6399:
URL: 
https://github.com/apache/datafusion-comet/issues/6399#issuecomment-5915354445

   Phase 3, planner, operators and shuffle: all 13 PRs have been reviewed 
against 1.0.0 and against Spark. Every reproducer was then run on 1.0.0 and 
1.1.0-rc1 builds, which covers Phase 4 for this area.
   
   Three more regressions that ship in 1.1.0 are confirmed, and #6402 has the 
details:
   
   - From #5667 and #5362: `explode`, `posexplode` and `explode_outer` of an 
array of structs with a boolean field return wrong booleans and NULLs after the 
first `spark.comet.batchSize` output rows of each input batch, with default 
configs. 3,000 rows of 12 elements gave wrong booleans in 22,824 of the 36,000 
output rows. The new chunking slices the output, and Arrow Java ignores the bit 
offset of the sliced boolean children, which is #6288. #6449, the `branch-1.1` 
backport of #6339, fixes every case on rc1, so the same backport covers this 
and #6424.
   - #5507, from #4870: the native `WindowGroupLimit` treats `-0.0` and `0.0` 
inside an array or struct sort key as different, so a `rank <= k` filter drops 
tied rows. 1.0.0 ran `WindowGroupLimit` in Spark.
   - From #5723: adaptive partial aggregation decides once per task, after 
100,000 rows, and never goes back. A task whose keys repeat after a mostly 
distinct start sends every later row through the shuffle, 10x the rows in the 
reproducer. The results are correct, and the only off switch is the testing 
config that the tuning guide documents.
   
   No regression was found in #5041, #5442, #5531, #5615, #5668, #5763 or 
#5803. #5421 and #5916 are deliberate trade-offs. #5421 moves partial `AVG`, 
`stddev`, `first`, `last` and a few others back to Spark when the final 
aggregate runs in Spark, mostly for anyone running without the Comet shuffle 
manager, because 1.0.0 could return NULL averages there. That's worth a line in 
the release notes. #5916 puts all of a map task's spill in one local directory, 
which #6010 tracks.
   
   One memory concern isn't verified. Since #5803, a native partial 
`collect_list` directly over `UNION ALL` or `coalesce` keeps every input batch 
it has read, including the columns it doesn't collect, and those aren't counted 
against the memory pool until the aggregate emits. I couldn't make that 
observable in a test on both builds.
   
   Why review missed these:
   
   - The explode tests exploded int arrays or pulled struct fields out into 
top-level columns, so a boolean below the top level never crossed back to the 
JVM from a sliced batch. #6288 names `LIMIT` and aggregates as triggers, not 
explode.
   - #5469 fixed the scalar case of the `WindowGroupLimit` bug and moved the 
nested case to #5507 as a documented incompatibility, instead of falling back 
for it by default.
   - #5723's tests covered uniformly high and uniformly low cardinality, but 
not a distinct prefix followed by repeats. Review asked for a supported off 
switch, and the PR documented the testing route instead.
   


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