grorge123 commented on PR #5526:
URL: 
https://github.com/apache/datafusion-comet/pull/5526#issuecomment-5928332788

   @sunchao @andygrove thank you both for the reviews. The branch is rebased 
onto today's main (`92b43cfaa`), and everything since the last push is in the 
last commit.
   
   **Your points**
   
   - **`LAST_WIN` (@sunchao):** `CometMapFromArrays` now runs the null guard 
before the `LAST_WIN` branch, so `allowIncompatible` can no longer hand a 
stateful child to `convert`. `map/map_from_arrays_last_win_opt_in.sql` and the 
gate test in `CometNullTypeCompositionSuite` both pin that order. #5846 landed 
in the meantime and fixes the evaluation order of `map_from_arrays` itself, so 
I rebased onto it and dropped my own version of that change. I kept one ANSI 
fixture from it, for values produced by a lambda the codegen dispatcher 
evaluates, since the new tests cover native values only.
   - **Published reasons (@andygrove):** `size`, `array_append`, `arrays_zip`, 
`collect_list` and `collect_set` now publish every reason they return, through 
shared `val`s. The gate test covers them, looking the aggregates up through 
`aggrSerdeMap`. It checks `size` under `spark.sql.legacy.sizeOfNull=false`, 
since Spark 3.x does not guard by default.
   
   **What else changed**
   
   Both points were instances of a wider pattern, so I went through the rest of 
the branch looking for more of the same:
   
   - `arrays_zip` had the evaluation-order problem that #5846 fixed for 
`map_from_arrays`: a single `AND` of null checks could evaluate a later 
argument on rows Spark's generated code skips, raising an ANSI error there. It 
now uses one `WHEN` per argument. There is an ANSI fixture for it, and it fails 
with the old guard. Spark's interpreted `ArraysZip.eval` does evaluate every 
argument, so where Spark would use it, `arrays_zip` over several arguments now 
falls back to Spark. Comet tracks those contexts while it serializes: 
`spark.sql.codegen.factoryMode=NO_CODEGEN`, beneath a non-leaf 
`CodegenFallback` expression (such as the `filter` that `array_compact` 
becomes, even with codegen on), the input of an imperative aggregate such as 
`collect_list` (but not its `FILTER`, which is generated code), and a generator 
that Spark leaves out of a whole-stage stage. In the same contexts the codegen 
dispatcher now compiles its kernel over the expression's `eval`, so a 
dispatched tree that contain
 s `arrays_zip` matches too, and `array_repeat` with a nullable count, whose 
`eval` skips the element when the count is NULL, is dispatched (on main too). 
The same holds for `trunc` / `date_trunc` with a non-literal format, whose 
`eval` returns NULL for an invalid format without evaluating the date (on main 
too, behind `allowIncompatible`). Each context has a fixture that fails without 
its fix.
   - I switched each gate off one at a time and compared results with Spark on 
the locked dependencies. Three gates turned out to be unnecessary, and I 
removed them:
     - The NullType `CASE` gate on `If` / `CaseWhen` / `Coalesce`. Native 
`CASE` handles NullType results on DataFusion 55.1.
     - The `collect_set` gate. The mismatches were only element order.
     - The `map_from_arrays` literal-beside-per-row gate. The native map 
expands the scalar.
   - `collect_list` keeps its NullType gate. Without a grouping key, the final 
merge reads partial state that never went through the planner's nullability 
cast. Plain nested types fail the same way on main (`collect_list(array(id))` 
over more than one partition), and NullType-bearing nested elements fail on 
main even in one partition.
   - A few guards were broader than needed and are narrowed:
     - `coalesce` no longer guards its last argument. A one-argument 
`coalesce`, which survives the optimizer only when `NullPropagation` and 
`SimplifyConditionals` are excluded, now serializes as its argument; it used to 
become a `CASE` with no `WHEN`, which native rejects (on main too).
     - `array_append` guards its item only when the array can be NULL.
   - `array(...)` over a NullType argument now goes through the dispatcher 
unless the argument is a literal. The old check used `foldable`, but a foldable 
argument that constant folding did not reduce (`element_at(array(NULL), 1)`) 
still arrives per row, and `make_array` returned one row for the whole batch.
   - `IF` is now serialized as a one-branch `CASE WHEN`. The native planner 
coerces `CASE` branches to one type but built `IfExpr` without that step, so 
two branches of one Spark type that differ in nested field nullability at 
native (native `named_struct` declares a field built from a literal 
non-nullable and one built from a column nullable) failed the projection's 
schema check whenever no row took the `THEN` branch. Since #6350 the `CASE` 
planner leaves a branch that already has the coerced type uncast, and evaluates 
`IF` and `CASE` through the same expression. A native `CASE` also names a 
merged struct's fields after its `ELSE` branch where Spark uses the first 
branch, which differ with case-insensitive analysis (`IF(c, s, 
named_struct('A', id))` with `s` a `struct<a>`); `array(...)`, `array_append` 
and the set ops then met two differently named structs. `IF`, `CASE WHEN` and 
`coalesce` now cast each branch of a struct-bearing type to the expression's 
own type. Comet's own casts 
 are subject to `spark.comet.expression.Cast.enabled`, which the expressions 
guide now notes.
   - The dispatcher now declares its output deep-nullable, like native 
constructors and Parquet columns. Two values of one Spark type used to reach 
native with different nested nullability depending on which side produced them, 
and kernels that compare their inputs' types (set ops, `array_remove` / 
`array_position`, `CASE`) rejected the pair. I found these with a sweep that 
pairs producers of different origin (native constructor, dispatcher, Parquet 
column, literal) of eight types inside about forty consumers that need one 
type, run on main and on this branch.
   - That sweep also turned up two problems that exist on main:
     - `arrays_overlap` with a literal array beside a column compared only the 
first row, or panicked (`arrays_overlap(array(id), array(1L))` over a multi-row 
batch). The kernel now expands the scalar side.
     - A dispatcher kernel whose generated code does not compile failed the 
query. Spark's own `ElementAt` codegen, for example, assigns an undeclared null 
flag when the index is an in-bounds literal and ANSI is off. The kernel now 
falls back to the expression's interpreted `eval`, as Spark's whole-stage 
codegen does.
   - A review round found a third one on main: a nested array literal whose 
inner arrays are all empty, such as a folded `array(cast(array() AS 
array<bigint>))`, was built one list level deeper than its declared type, so 
`array_union`, `arrays_overlap`, `IF`, `coalesce` and `=` against it failed on 
main. This branch newly reached the same bug through 
`arrays_overlap(transform(array(id), x -> array()), array(array()))`, whose 
dispatched side main declined. The literal builder now builds an empty child 
the way it builds a populated one, which keeps the declared depth and also 
fixes a deeper literal mixing the two (`array(array(array()), 
array(array(array(1))))`), which failed to build on main.
   - Set ops now cast both sides to the set op's own element type rather than 
each side's: with case-insensitive analysis Spark accepts 
`array(named_struct('a', x))` beside `array(named_struct('A', x))` without a 
cast, and the native kernel's type check compares field names, so `array_union` 
over such a pair failed (on main too, for typed structs; newly reachable here 
for NullType-bearing dispatched ones). Every side is cast, even one whose Spark 
type already matches: a native `CASE` names a merged struct's fields after its 
`ELSE` branch where Spark uses the first branch, so `array_union(array(IF(id < 
0, s, named_struct('A', id))), array(s))` failed too (also on main). A cast 
that only relabels nested names and nullability is now `Compatible`, so these 
casts stay native.
   - Native round-robin shuffle now falls back when a column it hashes contains 
a NullType. The native hasher has no NullType arm, so `REPARTITION(n)` over a 
dispatched `map(id, NULL)` failed (main declined that producer). Positional 
placement and columns past `maxHashColumns` are not hashed and stay native.
   - The non-AQE DPP rewrite now honours the broadcast gate, so the DPP 
subquery keeps reusing the join's Spark exchange.
   - The multi-batch shuffle test now reaches a second writer batch. It fails 
when the builder reset is removed; the previous version did not.
   
   **Coverage against main**
   
   I ran the same 6,901 queries on main (`ce455f32d`) and on this branch before 
the rebase: every query the composition sweep generates, with ANSI on and off, 
plus every query in the SQL fixtures.
   
   | main → this branch | queries |
   | --- | --- |
   | fallback → native | 1,482 |
   | error → native or fallback | 129 |
   | wrong answer → native | 3 |
   | native → fallback | 18 |
   
   Of the 18, 12 are deliberate and 6 come from a gate that is wider than it 
needs to be:
   - 8 are broadcast joins whose build side has a NullType directly under a 
struct. On main they run when the build side is one buffer and hang Arrow's 
appender when it is several. This branch skips coalescing those shapes so they 
cannot hang, but uncoalesced buffers cost each consuming task one IPC stream 
per buffer, which measured slower than Spark's broadcast, so the build side 
stays on Spark.
   - 4 are non-deterministic arguments under a null guard (`size`, 
`map_from_arrays`, `arrays_zip`). Main builds the same guards and serializes 
the argument twice. Its answers matched Spark on the sweep's data, but the two 
copies do not see the same rows, as in the `coalesce` case under the known 
issues below.
   - 6 are `collect_list` over a struct with a NullType field directly under it 
(`collect_list(named_struct('a', id, 'b', NULL))`), which main handles. The 
gate refuses any NullType under a struct to keep out the 16 shapes that fail on 
main, where that struct sits inside an array or another struct 
(`collect_list(array(named_struct('a', id, 'b', NULL)))`). Narrowing it to 
those shapes can be a follow-up.
   
   After each rebase I ran the branch side again. Apart from new fixture 
queries, the only changes come from upstream: `array_min` on signed zeros, and 
float noise in the probe's exact comparison of `stddev`, `corr` and `regr_*`.
   
   **Known issues left as they are**
   
   All of these are present on main, and none is made worse here:
   
   - Identical non-deterministic expressions that go through the codegen 
dispatcher share one kernel, so the second copy continues the first one's 
state. Two identical `coalesce(IF(monotonically_increasing_id() % 2 = 0, id, 
NULL), -1)` columns now reach this, where they used to evaluate the argument 
twice natively and were wrong in both columns.
   - Ungrouped `collect_list` over a nested input with non-nullable fields 
fails in the final merge (above).
   - `array_append` on Spark 3.x: Spark's generated code evaluates the item 
even for a NULL array, so an item that throws under ANSI raises in Spark and 
returns NULL natively.
   - `array_insert` (which `array_append` becomes on Spark 4.0+): Spark's 
`eval` returns NULL for a NULL array or position before it evaluates the item, 
while its generated code and the native kernel evaluate all three, so where 
Spark uses `eval` an item that throws raises natively only.
   - The non-AQE DPP rewrite follows the broadcast gate but not 
`spark.comet.exec.broadcastExchange.enabled`.
   - `array_join` with `allowIncompatible` and a nullable non-deterministic 
`nullReplacement` serializes the replacement twice, like the null guards above, 
but has no `NullGuard` check.
   - `array_join` with a nullable `nullReplacement` follows Spark's generated 
code, which skips the array when the replacement is NULL; Spark's interpreted 
`eval` evaluates the array first. So where Spark uses `eval` (for example the 
input of `collect_list`), an array argument that raises under ANSI raises in 
Spark and not natively. This is the same kind of `eval`/generated-code 
difference this PR handles for `arrays_zip`, `array_repeat` and `trunc`; I left 
`array_join` as it is on main rather than widen the PR further, and can open a 
follow-up issue for it.
   - The `map_from_arrays` compatibility note still says a NULL key is not 
detected. Native now rejects it, with a different message.
   
   **Testing**
   
   On the final commit, with the native library built in release mode:
   
   - Spark 4.1: `CometSqlFileTestSuite` (601), `CometNullTypeCompositionSuite` 
(28), `CometArrayExpressionSuite` (70), `CometExpressionSuite` (175), 
`CometTemporalExpressionSuite` (37), `CometMapExpressionSuite` (31), 
`CometAggregateSuite` (128), `CometJoinSuite` (56), `CometExecSuite` (154), 
`CometJsonExpressionSuite` (8), `CometCodegenSuite` (104), 
`CometCodegenSourceSuite` (67), `CometShuffleSuite` (48), 
`DisableAQECometShuffleSuite` (48), `CometNativeShuffleSuite` (58), 
`CometNativeCastSuite` (187), `UtilsSuite` (11), `GenerateDocsSuite` (4), the 
native unit tests for list literals and `CASE`, `cargo clippy`, `cargo fmt` and 
`spotless:check`.
   - Spark 3.5: `CometSqlFileTestSuite` (601), `CometNullTypeCompositionSuite` 
(28), `CometArrayExpressionSuite` (69), `CometExpressionSuite` (170), 
`CometAggregateSuite` (126), `CometCodegenSuite` (103) and 
`CometNativeShuffleSuite` (52).
   
   All passed, and each new regression test fails with its fix reverted. I also 
re-ran `CometNullTypeColumnsBenchmark` on the final commit. Each case runs the 
same plans as before, every arm returns the same rows, and every arm's speed 
relative to Spark is within 0.1x of the numbers I posted earlier, except the 
dispatcher-off arm of the 16K-row struct join, a run of about 45 ms that has 
ranged from 1.2x to 1.5x. I have not run Spark's own SQL suites locally. Since 
this touches the serde, the planner and the native shuffle writer, could a 
maintainer approve the workflow run and add `run-all-spark-profiles` and 
`run-spark-4.1-tests`?
   
   If the change still looks too broad, or its fallbacks cost more than it is 
worth, I'm happy to close this PR or split it into smaller ones.
   
   Assisted-by: Claude Code (claude-opus-5-5)
   


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