andygrove opened a new issue, #5676:
URL: https://github.com/apache/datafusion-comet/issues/5676

   ### What is the problem the feature request solves?
   
   Plain `posexplode` evaluates its array argument twice per batch. 
`posexplode_outer` does not — it already has the fix — so this is an 
inconsistency between two adjacent branches of the same planner match rather 
than a missing feature.
   
   `QueryPlanSerde` sends the generator's child as one expression, and 
`create_plan`'s `Explode` arm puts it into the pre-explode projection twice 
when `position` is set: once wrapped in `ListPositionsExpr` to produce `pos`, 
and once as the array column that `ExplodeExec` unnests.
   
   
https://github.com/apache/datafusion-comet/blob/main/native/core/src/execution/planner.rs#L2121-L2126
   
   `ProjectionExec` evaluates each expression in its list independently, so the 
child runs once per branch. When the child is an attribute that costs nothing, 
the second reference is an `Arc` clone. When it is a real expression it is real 
work, and when it is an expression Comet does not implement natively it is a 
JNI round-trip per batch, because codegen dispatch sends it back to the JVM as 
a `JvmScalarUdf`.
   
   The `outer` variant avoids this. Comet has to wrap the array in 
`ListEmptyToNullExpr` there to get Spark's outer semantics, and rather than 
re-running that wrapper per branch it materializes it in a pre-projection under 
a reserved name, so the upper projection references a column:
   
   
https://github.com/apache/datafusion-comet/blob/main/native/core/src/execution/planner.rs#L2063-L2092
   
   The comment above that block reasons about the wrapper only, and concludes 
that the non-outer cases reference the array once so need no pre-projection. 
That holds for `explode` and `explode_outer`. It does not hold for 
`posexplode`, where the raw child is the thing referenced twice.
   
   ### Reproduction
   
   `SELECT posexplode(split(s, 'x')) FROM t` against a one-column Parquet 
table, with `spark.comet.explain.native.enabled=true`. `split` is not native 
here, so it arrives as a `JvmScalarUdf`:
   
   ```
   CometExplodeExec, elapsed_compute=205.00µs
     ProjectionExec: expr=[list_positions(JvmScalarUdf(..., s@0)) as pos,
                           JvmScalarUdf(..., s@0)]
                     expr_0_eval_time=5.07ms, expr_1_eval_time=1.59ms
   ```
   
   Both projection expressions pay for the same `split`. The same query as 
`posexplode_outer` evaluates it once, in the pre-projection, and the upper 
projection is free:
   
   ```
   CometExplodeExec, elapsed_compute=2.50µs
     ProjectionExec: expr=[list_positions(__comet_explode_outer_...@1) as pos,
                           __comet_explode_outer_...@1 as ...]
                     expr_0_eval_time=2.00µs, expr_1_eval_time=41ns
       ProjectionExec: expr=[s@0 as s, list_empty_to_null(JvmScalarUdf(..., 
s@0)) as __comet_explode_outer_...]
                     expr_1_eval_time=4.15ms
   ```
   
   4.15ms of work once, against 5.07ms + 1.59ms for the same query without 
`outer`.
   
   This is invisible in the Spark plan, which shows one `CometProject` either 
way. It shows up as a `CometProject` under a `CometExplode` that is far slower 
than the explode above it — though note that projection is also where a 
non-native array expression legitimately runs, so a slow one there is not by 
itself evidence of this bug. The `expr_N_eval_time` metrics are what 
distinguish the two.
   
   ### Describe the potential solution
   
   Take the pre-projection whenever `position` is set, not only when `outer` 
is, wrapping with `ListEmptyToNullExpr` only in the outer case. The existing 
`(true, true)` branch is most of the code already; the match arms collapse to 
something like `(outer, true)` building the pre-projection over 
`raw_child_expr` or its wrapped form, and `(_, false)` unchanged.
   
   Two things worth getting right:
   
   - Skip the pre-projection when the child is already a `Column`, which is the 
common `posexplode(array_col)` shape. Otherwise every such query gains a 
`ProjectionExec` for no benefit.
   - The reserved output name is currently built as 
`format!("__comet_explode_outer_{}", child_field_name)`, which will read oddly 
once non-outer plans use it. Worth renaming while the code is being touched.
   
   No behavior change is expected, so `CometGenerateExecSuite` should pass 
untouched. The measurable claim is that `expr_1_eval_time` on the upper 
projection drops to near zero and total time for `posexplode` over a 
non-trivial array expression approaches `posexplode_outer` over the same one.
   
   ### Additional context
   
   Found while looking at explode performance more broadly. Related but 
separate: apache/datafusion-comet#5667 optimizes the unnesting kernels inside 
`ExplodeExec`, which is the operator above this projection and does not touch 
how the array expression is fed to it. Issue apache/datafusion-comet#5210 
tracks retiring `ListEmptyToNullExpr` in favor of upstream `unnest_outer`; if 
that lands first, the outer branch loses its wrapper but the double-reference 
in the `position` case remains, since it is about the raw child and not the 
wrapper.
   


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