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]