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

   ## Which issue does this PR close?
   
   Closes #6334.
   
   Found by the 1.1.0 regression audit (#6399) and tracked in #6402.
   
   ## Rationale for this change
   
   Spark only requires the two branches of `IF` to have the same type up to 
nullability, so it adds no cast when a struct field, map value or array element 
can be NULL in one branch and not in the other. Comet serialized each branch 
with its own Arrow type, and the native `IfExpr` declared the THEN branch's 
type as its output.
   
   - On a batch where no row took the THEN branch, `IfExpr` returned the ELSE 
array unchanged, and the projection failed with `column types must match schema 
types`.
   - On a batch that mixed the two, DataFusion's `CaseExpr` casts the ELSE rows 
to the THEN branch's type. That failed when the THEN branch was the less 
nullable one, with `Cannot cast nullable struct field 'x' to non-nullable 
field`, or `Found unmasked nulls for non-nullable StructArray field "value"` 
for a map.
   
   `CASE WHEN` doesn't have this problem, because `create_case_expr` in the 
planner casts the branches to a common type first.
   
   This became a 1.1.0 regression with #5452, which made folded map literals 
native. A folded `map('z', 0)` has `valueContainsNull = false` and a 
`MAP<STRING, INT>` column has `valueContainsNull = true`, so `IF(c, m, map('z', 
0))` and `IF(m IS NULL, map('z', 0), m)` fail on 1.1.0-rc1, while 1.0.0 ran 
them in Spark.
   
   ## What changes are included in this PR?
   
   The planner now builds `IF` the way it builds `CASE WHEN`. It takes the 
common type of the two branches from DataFusion's `type_union_coercion`, which 
is what `get_coerce_type_for_case_expression` folds over the `CASE WHEN` 
branches, and wraps a branch whose Arrow type differs from it in a Comet 
`Cast`. A branch that already has the common type is left alone, so an `IF` 
whose branches agree, which is the usual case, is planned as before.
   
   - The THEN branch goes first, so the common type takes its field names from 
the THEN branch, as Spark's `If.dataType` does. Nullability is combined the 
same way in either order.
   - The cast uses the same options as the `CASE WHEN` casts, including the UTC 
timezone from #6347.
   - The branches share a Spark type, so the casts don't change any values. 
They make a nested field nullable, or relabel a field name or a timestamp's 
timezone. The native `Cast` already handles this for structs (field by field), 
lists (rebuilt with the target element field) and maps (the relabel path in 
`cast_map_to_map`). `CASE WHEN` relies on the same casts.
   
   I did this in the native planner rather than in `CometIf`, because that's 
where `CASE WHEN` does it and because it compares the Arrow types the branches 
actually produce. A JVM-side cast to Spark's `If.dataType` would miss a branch 
whose native type differs from its Spark type.
   
   This touches the same code as #6350, which replaces `create_case_expr` with 
`create_case_when` and changes `IfExpr`. This change doesn't add to #6350's 
conflict with main, which is the one with #6347 in `create_case_expr`, but 
whichever of the two lands second may need a rebase. Once #6350 is in, `IF` 
could go through `create_case_when` instead.
   
   ## How are these changes tested?
   
   Every new query runs with `checkSparkAnswerAndOperator`, and each table is 
written by three INSERTs. Each INSERT writes its own files and a batch never 
spans files, so every query sees batches in which every row takes the THEN 
branch, batches in which every row takes the ELSE branch, and batches that mix 
the two.
   
   - `if_nested_nullability.sql` covers the constructor path, since the SQL 
harness disables `ConstantFolding`. It runs the issue's struct and map queries 
in both branch orders, `IF(q, m, map('z', 0))` and `IF(m IS NULL, map('z', 0), 
m)` over a map column, a struct column against a struct constructor, and a 
struct and an array inside a map value. It also runs `IF(q, array(i), 
array(0))` and the `CASE WHEN` form, which already passed.
   - A new `CometMapExpressionSuite` test covers the folded map literal path 
with `ConstantFolding` on. It runs the audit's two queries and `IF(c, map('z', 
0), m)`.
   
   Without the fix, all 12 `IF` queries whose branches differ in nullability 
fail (9 in the SQL file and 3 in the Scala test). The two queries that already 
passed still pass. With the fix, all of them pass.
   
   These pass:
   
   - Spark 4.1 (the default profile), 625 tests: `./mvnw test -Dtest=none 
-Dsuites="org.apache.comet.expressions.conditional.CometIfSuite,org.apache.comet.expressions.conditional.CometCaseWhenSuite,org.apache.comet.expressions.conditional.CometCoalesceSuite,org.apache.comet.CometMapExpressionSuite,org.apache.comet.CometSqlFileTestSuite"`
   - Spark 3.4, 3.5 and 4.0, 17 tests each: `CometIfSuite`, the new 
`CometMapExpressionSuite` test and the `expressions/conditional/` SQL files.
   - `cargo test -p datafusion-comet --lib planner` and `cargo clippy 
--all-targets --workspace -- -D warnings`
   
   Before #6347 merged, I also ran `CometIfSuite`, `CometCaseWhenSuite`, the 
new test and the conditional SQL files on a merge of this branch with #6350. 
All 22 passed, including #6350's `case_when_*.sql` files.
   


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