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

   Auto-insert CometSparkToColumnarExec at an unsupported 
FileSourceScan/BatchScan on a broadcast join's build side when the probe side 
is natively scannable, so the BroadcastExchange and the join run natively over 
the large probe. Gated by 
spark.comet.sparkToColumnar.broadcastBuildSide.enabled (default true) together 
with the Comet broadcast-exchange config. Guard the AQE DPP rule (convertSAB) 
against wrapping a non-native DPP subquery build in a 
CometBroadcastExchangeExec.
   
   
   ## Which issue does this PR close?
   
   Part of #6008.
   
   ## Rationale for this change
   
   An unsupported leaf scan (for example a Text file source, `[COMET: 
Unsupported file format Text]`)
   on the **build side of a broadcast join** cascades: the build branch (scan → 
transforms →
   `BroadcastExchange`) stays on Spark, so the exchange never becomes a 
`CometBroadcastExchangeExec`,
   so the `BroadcastHashJoin` can't go native — even when the large probe side 
is a fully native
   scan+filter. A tiny lookup table read from an unsupported format 
disqualifies native execution of
   the entire join over the big stream.
   
   Comet already has a row→Arrow bridge (`CometSparkToColumnarExec`, delivered 
as a `CometScanWrapper`
   that is a `CometNativeExec`), but it is off by default and its 
`supportedOperatorList` excludes file
   scans, so this path was never hit. This PR auto-inserts that bridge at an 
unsupported build-side
   scan — but only when doing so actually lets the join run natively — so the 
bounded, small build side
   is adapted to Arrow and the join over the large probe input is accelerated.
   
   ## What changes are included in this PR?
   
   - **New config** `spark.comet.sparkToColumnar.broadcastBuildSide.enabled` 
(default `true`,
     `CometConf.scala`). Independent of the general 
`spark.comet.sparkToColumnar.enabled` opt-in.
   
   - **Auto-bridge on broadcast build sides** (`CometExecRule.scala`): a 
pre-pass
     (`tagBroadcastBuildSideLeaves`) tags the build-side 
`FileSourceScanExec`/`BatchScanExec` of a
     `BroadcastHashJoinExec`/`BroadcastNestedLoopJoinExec`, and 
`shouldApplySparkToColumnar` then
     bridges a tagged scan — only in the unsupported-format arms, so the 
per-format
     `spark.comet.convert.{csv,json,parquet}.enabled` opt-outs still take 
precedence.
     The bridge is applied only when:
     - the feature config **and** `spark.comet.exec.broadcastExchange.enabled` 
are both on
       (otherwise the broadcast can't go native, so the bridge would be 
pointless), and
     - the **probe side is already natively scannable** (`hasOnlyNativeScans`) 
— i.e. the join can
       actually become a fully native `CometBroadcastHashJoinExec`. Bridging a 
build side under a join
       that stays on Spark is wasteful and would change the build broadcast's 
subtree, breaking Spark's
       dynamic-partition-pruning (DPP) broadcast reuse.
     The tag descent walks the dimension's own filters/projects/aggregations 
but stops at a nested
     join, so a large streamed input of a nested join inside the build side is 
not bridged.
   
   - **DPP guard** (`CometPlanAdaptiveDynamicPruningFilters.convertSAB`): when 
the matched broadcast
     join is Comet, only build a `CometBroadcastExchangeExec` for the reused 
DPP subquery if that
     subquery's own build is native (`isNativeBuildSide`); otherwise fall back 
to a Spark
     `BroadcastExchangeExec`. This prevents an AQE+DPP crash
     (`Comet execution only takes Arrow Arrays, but got ... 
OffHeapColumnVector`) when the bridged
     dimension's separate DPP build copy is row-based. Mirrors the existing 
non-AQE guard in
     `CometExecRule.rewriteInSubqueryPlan`.
   
   ## How are these changes tested?
   
   - New unit tests in `CometExecSuite` (`#6008`), covering: a Text build side 
going native (AQE on
     and off); the broadcast-exchange-disabled gate; an unsupported scan below 
a build-side aggregate
     (and that the shuffle query stage is not bridged); the per-format opt-out 
still winning; a DSv2
     (`BatchScanExec`) build side; `BuildLeft` and `BroadcastNestedLoopJoin` 
variants; a non-native
     probe leaving the build side un-bridged; and an AQE+DPP regression test 
(text dimension + a
     partitioned Parquet fact) that pins the DPP guard. All assert 
result-equivalence to Spark via
     `checkSparkAnswer`/`checkSparkAnswerAndOperator` plus the expected plan 
shape.
   - `CometExecSuite` passes on `spark-4.1`; the project builds on `spark-4.0` 
and `spark-4.1`.
   - The Spark SQL suite via `dev/local-ci.sh spark sql_core-1` is green, 
including the
     `DynamicPartitionPruning` V1/V2 suites that guard broadcast/DPP reuse (an 
earlier eager version of
     this change regressed those; the probe-native gate + DPP guard fix it).
   
   


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