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

   ## Which issue does this PR close?
   
   Closes #5710. Part of #5572. Supersedes #5714.
   
   ## Rationale for this change
   
   A typed `Dataset` operation such as `ds.map(f)` plans a 
`DeserializeToObject` / `MapElements` / `SerializeFromObject` island that has 
to run in Spark, because it passes JVM objects between its operators. Comet 
only takes over again at the next shuffle, so whatever sits between the island 
and that shuffle stays on Spark too. For `ds.map(f).groupBy(k).agg(...)` that 
is the partial aggregate, which is expensive when there are many groups. A 
broadcast join on the probe side is in the same position.
   
   #5714 tried to fuse the island into one projection in the JVM codegen 
dispatcher. Since #6542 the dispatcher only calls into Spark's own classes, so 
that approach is blocked. When I measured the alternative in this PR against 
it, the conversion got 93-100% of the fused speedup in the cases where either 
one pays off. It is also much simpler, and it covers every typed operation 
rather than only `map`.
   
   ## What changes are included in this PR?
   
   Every typed operation ends in `SerializeFromObjectExec`, whose output is 
ordinary rows. With the new `spark.comet.convert.typedDataset.enabled`, 
`CometExecRule` puts a `CometSparkToColumnarExec` above it, so the operators 
above the typed operation can run natively. The operation itself, including the 
user function, runs in Spark exactly as before. This covers `map`, `flatMap`, 
`mapPartitions`, `groupByKey(...).mapGroups`, `cogroup`, and a Dataset built 
from an RDD of objects. If nothing native consumes the output, 
`EliminateRedundantTransitions` removes the conversion again, so a typed 
operation at the top of a plan is unchanged.
   
   Spark inserts no columnar transitions below a `RowToColumnarTransition`, and 
`CometSparkToColumnarExec` is one. Above a leaf that does not matter. Here the 
typed operation's own operators sit below the conversion, and without 
transitions they read their Comet child through `CometExec.doExecute`, which is 
Spark's interpreted columnar-to-row path. That gives the right answer slowly, 
and in the benchmark it made whole queries 1.5-2.4x slower. The rule therefore 
applies Spark's own `ApplyColumnarRulesAndInsertTransitions` to the subtree, as 
`CometRule.buildPreview` already does, and `EliminateRedundantTransitions` 
swaps in Comet's columnar-to-row as usual. Spark's rule leaves existing 
transitions alone, which matters because `CometExecRule` runs over the same 
plan twice under AQE.
   
   That second pass used to tag the inserted `ColumnarToRowExec` with 
"ColumnarToRow is not supported", so `CometExecRule` no longer reports a 
`ColumnarToRowTransition` as an operator it failed to convert. A column type 
that Spark-to-Comet conversion does not support, such as `array<int>`, keeps 
the operators above the typed operation on Spark and records a fallback reason 
that names the column.
   
   It is off by default because whether it pays depends on the work above the 
typed operation, which the planner cannot see. Here is 
`CometTypedDatasetBenchmark` on an M3 Max with 4Mi rows, `local[1]`, AQE off 
and one shuffle partition. Times are best of the iterations in ms, and the last 
column repeats the default Comet arm as a noise check.
   
   | case                               | Spark | Comet (default) | Comet, 
converted | default repeat |
   | ---------------------------------- | ----- | --------------- | 
---------------- | -------------- |
   | map -> group by 100 keys           | 249   | 248             | 313         
     | 221            |
   | map -> group by 1M keys            | 1244  | 1079            | **484**     
     | 940            |
   | map -> filter -> group by 100 keys | 180   | 153             | 319         
     | 153            |
   | map -> filter -> group by 1M keys  | 896   | 633             | **451**     
     | 636            |
   | map -> group by long key, 1M keys  | 811   | 699             | **421**     
     | 686            |
   | map, 4 columns -> group by 1M keys | 1359  | 1101            | **682**     
     | 1103           |
   | mapPartitions -> group by 1M keys  | 1289  | 1148            | **738**     
     | 1136           |
   
   With an aggregate over many groups above the typed operation, it is 1.4-2.2x 
faster than today's default. With a cheap aggregate it is slower, 0.5x after a 
selective filter, because Spark compiles the typed operation, the filter and a 
small aggregate into one loop, while the conversion writes every row to Arrow 
first.
   
   Two behaviors worth knowing when it is on:
   
   - An exception thrown by the user function reaches the driver as 
`CometNativeException: C Data interface error: 
java.lang.IllegalArgumentException: ...` rather than as the original exception. 
Every Spark-to-Arrow input wraps errors this way today (#6234), and the open 
#6243 fixes it for all of them.
   - The shuffle above the typed operation, which today is Comet's columnar 
shuffle, becomes a native shuffle. That is the choice Comet makes for any 
native child. The native partitioner's divergence on decimals wider than 18 
digits (#5994) applies to it like any other native shuffle.
   
   The user guide's operator page and its section on Spark-to-Comet conversion 
types describe the new config. The contributor guide's paragraph on typed 
operators suggested the dispatcher fuse, so it now describes the conversion and 
why the fuse was dropped.
   
   ## How are these changes tested?
   
   New `CometTypedDatasetSuite` has 10 tests, registered in both PR workflows. 
Most of them check the answer against Spark, that every operator above the 
typed operation is native, and that the conversion sits directly on 
`SerializeFromObjectExec`. They cover:
   
   - `map` followed by an aggregate with AQE on and off, and `flatMap`, 
`mapPartitions`, `mapGroups`, `cogroup`, and an RDD of objects.
   - A broadcast join above the typed operation running as 
`CometBroadcastHashJoinExec`.
   - Decimal, `Option`, nested struct and `array<string>` fields, and an 
`array<int>` column that declines with its fallback reason.
   - A typed operation at the top of the plan, where nothing is converted.
   - The user function running exactly once per row.
   - An exception thrown past the first batch failing the query with the 
original message.
   - The feature being off by default.
   
   Every converted plan is also checked for row operators reading a columnar 
child without a transition, and for spurious `ColumnarToRow` fallback reasons. 
Dropping the transition insertion fails 5 tests, and dropping the 
`ColumnarToRowTransition` change fails 4.
   
   The suite passes locally on Spark 3.4, 3.5, 4.0, 4.1 and 4.2. On 4.1 these 
also pass: `CometExecSuite`, `CometExecRuleSuite`, 
`EliminateRedundantTransitionsSuite`, 
`RevertNativeForTransitionHeavyStagesSuite`, `CometInMemoryCacheSuite`, 
`CometRangeExecSuite`, `CometJoinSuite`, and both TPC-DS plan stability suites 
with no approved plan changes. Semantic scalafix (3.4), scalastyle, spotless 
and prettier are clean.
   
   I also ran Spark's own typed `Dataset` suites from the 4.1.3 tests jar, with 
Comet enabled through system properties: `DatasetSuite`, 
`DatasetPrimitiveSuite`, `DatasetAggregatorSuite`, `DatasetOptimizationSuite`, 
`DatasetCacheSuite` and `DatasetSerializerRegistratorSuite`, 281 tests. With 
the conversion off and on, the same 278 pass and the same 3 fail: the two 
TIME-type tests, which need Spark's test-only confs, and `groupBy.as`, which 
asserts on Spark's own plan nodes. A query listener showed that 22 of the 
suites' queries ran with the conversion when it was on.
   


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