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]
