andygrove commented on code in PR #6786:
URL: https://github.com/apache/datafusion-comet/pull/6786#discussion_r4231058284
##########
spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:
##########
@@ -1518,6 +1535,126 @@ case class CometExecRule(session: SparkSession,
queryStagePrep: Boolean = false)
}
}
+ /**
+ * An unsupported file-source scan is bridged to Arrow when the general
sparkToColumnar opt-in
+ * covers it, or when it feeds a broadcast join's build side (see issue
#6008). This is only
+ * reached from the unsupported-format arms, so the per-format
`COMET_CONVERT_FROM_*` opt-outs
+ * for CSV/JSON/Parquet still take precedence.
+ */
+ private def shouldBridgeUnsupportedScan(conf: SQLConf, op: SparkPlan):
Boolean =
+ isSparkToArrowEnabled(conf, op) || isBroadcastBuildSideLeaf(op)
+
+ /**
+ * True when the build-side auto-bridge should run: its own config is on and
Comet broadcast
+ * exchange conversion is enabled. Without the latter the
`BroadcastExchangeExec` cannot become
+ * a `CometBroadcastExchangeExec`, so bridging the build side would only add
a row to Arrow copy
+ * for no native join. This does not also check the join serdes
+ * (COMET_EXEC_BROADCAST_HASH_JOIN_ENABLED /
COMET_EXEC_BROADCAST_NESTED_LOOP_JOIN_ENABLED): if
+ * the exchange converts but the join does not, the broadcast arm returns
the original plan and
+ * EliminateRedundantTransitions removes the now-unused bridge, so a
leftover copy never reaches
Review Comment:
I don't think the conversion always goes away when the join stays on Spark.
`EliminateRedundantTransitions` only removes it when it sits directly under the
exchange. With a filter or project in between, which is the usual shape for a
dimension, it stays. With default settings, a `FULL OUTER` nested loop join
against a Parquet probe ends up with `BroadcastExchange <- CometColumnarToRow
<- CometFilter <- CometSparkRowToColumnar <- FileScan text`. So does a `LEFT
OUTER` one that builds its left side. Today that build side is a plain Spark
filter over the scan.
The same thing happens whenever the hash join declines, for example with
`spark.comet.exec.broadcastHashJoin.enabled=false`. With DPP on top, the main
broadcast no longer matches the DPP subquery's, so the dimension is scanned and
broadcast twice. That's the cost the probe-side check was meant to avoid, but
the check only looks at scans. Could the rule restore the original build
subtree when the join doesn't convert? That would also fix the 3.4 case.
##########
spark/src/main/scala/org/apache/comet/rules/CometPlanAdaptiveDynamicPruningFilters.scala:
##########
@@ -351,6 +360,26 @@ case object CometPlanAdaptiveDynamicPruningFilters
}
}
+ /**
+ * True when `plan` (a DPP subquery's build, taken from an
AdaptiveSparkPlanExec.executedPlan)
+ * produces native Comet columnar output, so it is safe as a
CometBroadcastExchangeExec child.
+ * Unwraps AQE stage wrappers and any columnar->row transition (mirrors the
stripping in
+ * CometExecRule.rewriteInSubqueryPlan), then checks for a CometNativeExec.
See issue #6008.
+ */
+ private def isNativeBuildSide(plan: SparkPlan): Boolean = {
+ val stripped = stripAQEPlan(plan) match {
+ case c2r: CometNativeColumnarToRowExec => c2r.child
+ case c2r: CometColumnarToRowExec => c2r.child
+ case WholeStageCodegenExec(c2r: CometColumnarToRowExec) =>
+ c2r.child match {
+ case InputAdapter(child) => child
+ case other => other
+ }
+ case other => other
+ }
+ stripped.isInstanceOf[CometNativeExec]
Review Comment:
This check runs even with the new config off, and it rejects Comet plans
that the join's own broadcast accepts. `CometUnionExec`, `CometCoalesceExec`
and `CometTakeOrderedAndProjectExec` extend `CometExec`, not `CometNativeExec`.
Take an AQE DPP query whose dimension is a `UNION ALL` of two filtered Parquet
tables. Before this change, the DPP subquery is a `CometSubqueryBroadcast` over
a `ReusedExchange` of the join's `CometBroadcastExchange`. With it, the
subquery builds its own Spark `BroadcastExchange` over
`CometColumnarToRow(CometUnion(...))`, so the dimension runs twice. I saw this
on 3.5, 4.1 and 4.2, and putting the old line back restores the reuse.
Could this check for any Comet columnar plan instead, for example
`stripped.isInstanceOf[CometPlan] && stripped.supportsColumnar`? With that
change the reuse comes back, the new text test still takes the Spark path, and
the existing DPP tests pass. A union test that asserts the `ReusedExchangeExec`
would lock this in. The non-AQE guard in `rewriteInSubqueryPlan` has the same
gap on main today. I filed #6815 for that one.
##########
spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala:
##########
@@ -301,6 +301,366 @@ class CometExecSuite extends CometTestBase {
}
}
+ test("broadcast build side with unsupported text source goes native
(#6008)") {
Review Comment:
All the new tests read text, but the catch-all branches send every other
file format and V2 scan through the conversion too. That covers ORC, Avro, XML,
Iceberg tables the native reader declines, Iceberg metadata tables such as
`t.snapshots`, and JDBC and other V2 connectors. Nothing in Comet's suites
reads ORC through `CometSparkToColumnarExec` today.
I ran an ORC dimension with nulls in every primitive type, a struct,
`ARRAY<STRING>`, `MAP<STRING,STRING>`, `TIMESTAMP` and `TIMESTAMP_NTZ`, in an
`America/Los_Angeles` session. I used V1 and V2, and the vectorized, nested
vectorized and row readers. I also ran Iceberg ORC, Iceberg merge-on-read with
native scans off, and `snapshots`. Everything matched Spark on 3.4, 3.5 and
4.1, and ORC also matched on 4.2. Could you add an ORC build-side test like
that, plus one for an Iceberg table the native reader declines? I'm happy to
share the probe.
##########
spark/src/main/scala/org/apache/comet/CometConf.scala:
##########
@@ -1166,6 +1166,20 @@ object CometConf extends ShimCometConf {
.toSequence
.createWithDefault(Nil)
+ val COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED: ConfigEntry[Boolean] =
+ conf("spark.comet.sparkToColumnar.broadcastBuildSide.enabled")
Review Comment:
#6602 moved the per-source conversion switches under
`spark.comet.convert.*`, and `config_conventions.md` lists `convert` as the
category for these. Could this be
`spark.comet.convert.broadcastBuildSide.enabled` while it's still easy to
rename? The doc text says it applies to "a leaf that Comet cannot scan
natively". It only covers file and V2 scans, though, and only when every
probe-side scan is a Comet scan. Could it say that, and also mention that with
DPP the dimension is read twice?
##########
spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:
##########
@@ -1518,6 +1535,126 @@ case class CometExecRule(session: SparkSession,
queryStagePrep: Boolean = false)
}
}
+ /**
+ * An unsupported file-source scan is bridged to Arrow when the general
sparkToColumnar opt-in
+ * covers it, or when it feeds a broadcast join's build side (see issue
#6008). This is only
+ * reached from the unsupported-format arms, so the per-format
`COMET_CONVERT_FROM_*` opt-outs
+ * for CSV/JSON/Parquet still take precedence.
+ */
+ private def shouldBridgeUnsupportedScan(conf: SQLConf, op: SparkPlan):
Boolean =
+ isSparkToArrowEnabled(conf, op) || isBroadcastBuildSideLeaf(op)
+
+ /**
+ * True when the build-side auto-bridge should run: its own config is on and
Comet broadcast
+ * exchange conversion is enabled. Without the latter the
`BroadcastExchangeExec` cannot become
+ * a `CometBroadcastExchangeExec`, so bridging the build side would only add
a row to Arrow copy
+ * for no native join. This does not also check the join serdes
+ * (COMET_EXEC_BROADCAST_HASH_JOIN_ENABLED /
COMET_EXEC_BROADCAST_NESTED_LOOP_JOIN_ENABLED): if
+ * the exchange converts but the join does not, the broadcast arm returns
the original plan and
+ * EliminateRedundantTransitions removes the now-unused bridge, so a
leftover copy never reaches
+ * runtime. See issue #6008.
+ */
+ private def broadcastBuildSideBridgeEnabled: Boolean =
+ CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.get(conf) &&
+ CometConf.COMET_EXEC_BROADCAST_EXCHANGE_ENABLED.get(conf)
+
+ /**
+ * True when `op` is a scan on a broadcast join's build side and build-side
bridging is enabled.
+ * `tagBroadcastBuildSideLeaves` only sets the tag when bridging is enabled,
but the tag lives
+ * on the physical node and can survive onto a reused instance, so re-check
here: a stale tag
+ * must not bridge a scan once the feature is turned off. A tagged scan is
bridged to Arrow even
+ * when the general sparkToColumnar path is off, letting the build branch
and the join run
+ * natively. See issue #6008.
+ */
+ private def isBroadcastBuildSideLeaf(op: SparkPlan): Boolean =
+ broadcastBuildSideBridgeEnabled &&
+ op.getTagValue(CometExecRule.BROADCAST_BUILD_SIDE_TAG).isDefined
+
+ /**
+ * Tag the file-source scan leaves on a broadcast join's build side so
+ * shouldApplySparkToColumnar can bridge an unsupported build-side scan
(e.g. a Text scan) to
+ * Arrow, letting the whole join run natively. The build side is usually
small (auto-broadcasts
+ * are capped by the broadcast threshold), so the row to Arrow copy is
usually cheap - but that
+ * is a heuristic, not a bound: an explicit BROADCAST hint can force a
larger build side, and a
+ * selective filter above the scan means the bridge copies the pre-filter
scan output. See issue
+ * #6008.
+ *
+ * The bridge is applied ONLY when the probe side is already natively
scannable, i.e. the join
+ * can actually become a fully native CometBroadcastHashJoinExec. Bridging a
build side under a
+ * join that stays on Spark is pure overhead and, worse, it rewrites the
build broadcast's
+ * subtree to Comet and breaks Spark's DPP broadcast reuse (the DPP subquery
keeps an unbridged
+ * copy of the same scan, so the two exchanges no longer share a canonical
form). CometScanRule
+ * has already run by this point, so a native probe scan is a Comet scan and
an unsupported
+ * probe scan is still a plain Spark scan - see `hasOnlyNativeScans`.
+ *
+ * Only `FileSourceScanExec` / `BatchScanExec` are tagged - the scan types
the per-format arms
+ * of shouldApplySparkToColumnar handle. A natively scannable file (e.g.
Parquet) is normally
+ * already a `CometScanExec` here; if native scan conversion did not take,
its
+ * `FileSourceScanExec` is still tagged, but the tag is harmless because the
per-format arm
+ * decides that scan before isBroadcastBuildSideLeaf is ever consulted. DPP
subquery broadcasts
+ * live in expressions, not `children`, so `foreach` never reaches them -
that scope is
+ * intentional.
+ */
+ private def tagBroadcastBuildSideLeaves(plan: SparkPlan): Unit = {
+ if (!broadcastBuildSideBridgeEnabled) {
+ return
+ }
+ plan.foreach {
+ case j: BroadcastHashJoinExec => tagBuildIfProbeNative(j.buildSide,
j.left, j.right)
+ case j: BroadcastNestedLoopJoinExec =>
tagBuildIfProbeNative(j.buildSide, j.left, j.right)
+ case _ =>
+ }
+ }
+
+ /** Tag the build side's file-source scans, but only if the probe side is
natively scannable. */
+ private def tagBuildIfProbeNative(
+ buildSide: BuildSide,
+ left: SparkPlan,
+ right: SparkPlan): Unit = {
+ val (buildPlan, probePlan) = buildSide match {
+ case BuildLeft => (left, right)
+ case BuildRight => (right, left)
+ }
+ if (hasOnlyNativeScans(probePlan)) {
Review Comment:
On Spark 3.4 with AQE, this turns off dynamic partition pruning for a native
Iceberg fact joined to a text dimension, and the join still runs in Spark.
`CometSpark34AqeDppFallbackRule` puts `SKIP_COMET_BROADCAST_TAG` on the
build-side `BroadcastExchangeExec`. That keeps the broadcast on Spark so that
Spark's `PlanAdaptiveDynamicPruningFilters` can match it with `sameResult`. The
pre-pass doesn't check that tag, so it converts the text scan under the Spark
broadcast anyway. The broadcast's subtree now has Comet nodes in it and the DPP
subquery's copy doesn't, so the match fails and DPP becomes `true`.
I reproduced this on 3.4.3. The fact is partitioned by `fk`, and the
dimension is filtered on `CAST(value AS INT) % 3 = 0` so Spark can't infer that
filter onto the fact. With the new config on, the fact scan returns 1000 rows.
With it off, it returns 400 and the broadcast is reused. Could
`tagBuildIfProbeNative` skip a join whose build child is a
`BroadcastExchangeExec` with that tag? I tried this locally. DPP comes back on
3.4 and the new tests still pass.
##########
spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala:
##########
@@ -301,6 +301,366 @@ class CometExecSuite extends CometTestBase {
}
}
+ test("broadcast build side with unsupported text source goes native
(#6008)") {
+ // A tiny lookup table read from a Text file (which Comet cannot scan
natively) sits on the
+ // build side of a broadcast join over a large native probe. Without the
auto-bridge the whole
+ // join stays on Spark. With
spark.comet.sparkToColumnar.broadcastBuildSide.enabled (default on)
+ // the text leaf is bridged to Arrow so the broadcast and the join run
natively.
+ assert(
+
CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.defaultValue.contains(true))
+ withTempDir { dir =>
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ spark
+ .range(0, 5)
+ .selectExpr("CAST(id AS STRING) AS value")
+ .write
+ .text(allowPath)
+ }
+
+ // Cover both the AQE and non-AQE planning paths - the tag pre-pass runs
per stage under AQE.
+ for (aqe <- Seq(false, true)) {
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe.toString) {
+ spark.read.parquet(probePath).createOrReplaceTempView("probe_6008")
+ spark.read.text(allowPath).createOrReplaceTempView("allow_6008")
+ val query =
+ "SELECT /*+ BROADCAST(a) */ p.v FROM probe_6008 p JOIN allow_6008
a ON p.key = a.value"
+
+ // Default: the text build side is bridged to Arrow and the join
goes native. The plan is
+ // fully native (the bridge counts as a native leaf), so use
checkSparkAnswerAndOperator -
+ // checkSparkAnswer would fail under COMET_STRICT_TESTING on a
fully-native plan.
+ val (_, cometPlan) = checkSparkAnswerAndOperator(sql(query))
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s
}.nonEmpty,
+ s"aqe=$aqe: expected a CometSparkToColumnarExec on the text build
side:\n" +
+ cometPlan.treeString)
+ assert(
+ collect(cometPlan) { case b: CometBroadcastExchangeExec => b
}.nonEmpty,
+ s"aqe=$aqe: expected a
CometBroadcastExchangeExec:\n${cometPlan.treeString}")
+ assert(
+ collect(cometPlan) { case j: CometBroadcastHashJoinExec => j
}.nonEmpty,
+ s"aqe=$aqe: expected a
CometBroadcastHashJoinExec:\n${cometPlan.treeString}")
+
+ // Off: the unsupported text leaf blocks native execution of the
broadcast and the join.
+ withSQLConf(
+ CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.key ->
"false") {
+ val (_, cometPlanOff) = checkSparkAnswer(sql(query))
+ assert(
+ collect(cometPlanOff) { case s: CometSparkToColumnarExec => s
}.isEmpty,
+ s"aqe=$aqe: expected no CometSparkToColumnarExec when
disabled:\n" +
+ cometPlanOff.treeString)
+ assert(
+ collect(cometPlanOff) { case b: CometBroadcastExchangeExec => b
}.isEmpty,
+ s"aqe=$aqe: expected no CometBroadcastExchangeExec when
disabled:\n" +
+ cometPlanOff.treeString)
+ assert(
+ collect(cometPlanOff) { case j: CometBroadcastHashJoinExec => j
}.isEmpty,
+ s"aqe=$aqe: expected no CometBroadcastHashJoinExec when
disabled:\n" +
+ cometPlanOff.treeString)
+ }
+ }
+ }
+ }
+ }
+
+ test("broadcast build side bridge is gated on broadcast-exchange conversion
(#6008)") {
+ // The bridge is pointless when Comet broadcast exchange conversion is
off: the broadcast (and
+ // so the join) cannot go native. With
COMET_EXEC_BROADCAST_EXCHANGE_ENABLED=false the text
+ // build side must not be bridged even though broadcastBuildSide is on.
+ withTempDir { dir =>
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ spark.range(0, 5).selectExpr("CAST(id AS STRING) AS
value").write.text(allowPath)
+ }
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.key ->
"true",
+ CometConf.COMET_EXEC_BROADCAST_EXCHANGE_ENABLED.key -> "false") {
+
spark.read.parquet(probePath).createOrReplaceTempView("probe_gate_6008")
+ spark.read.text(allowPath).createOrReplaceTempView("allow_gate_6008")
+ val (_, cometPlan) = checkSparkAnswer(sql("""
+ |SELECT /*+ BROADCAST(a) */ p.v
+ |FROM probe_gate_6008 p JOIN allow_gate_6008 a ON p.key = a.value
+ |""".stripMargin))
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s }.isEmpty,
+ "with broadcast-exchange conversion off the build side must not be
bridged:\n" +
+ cometPlan.treeString)
+ }
+ }
+ }
+
+ test("broadcast build side scan below an aggregate is bridged, shuffle stage
is not (#6008)") {
+ // A file scan can sit BELOW build-side operators. Broadcasting an
aggregated Text table puts a
+ // Text FileSourceScanExec under a partial aggregate + shuffle on the
build side. The tag
+ // pre-pass walks the whole broadcast subtree, so the deep text scan must
be bridged to Arrow,
+ // yet no query stage (e.g. the build-side shuffle under AQE) may be
wrapped in the bridge.
+ withTempDir { dir =>
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark.range(0, 20).selectExpr("CAST(id % 5 AS STRING) AS
value").write.text(allowPath)
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id % 5 AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ }
+ spark.read.text(allowPath).createOrReplaceTempView("allow_agg_6008")
+ spark.read.parquet(probePath).createOrReplaceTempView("probe_agg_6008")
+ val query = """
+ |SELECT /*+ BROADCAST(d) */ p.v, d.n
+ |FROM probe_agg_6008 p
+ |JOIN (SELECT value AS k, COUNT(*) AS n FROM allow_agg_6008 GROUP BY
value) d
+ | ON p.key = d.k
+ |""".stripMargin
+ for (aqe <- Seq(false, true)) {
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe.toString,
+ SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> (10 * 1024 *
1024).toString) {
+ // The plan is fully native (the bridged text scan under the
aggregate is a native leaf),
+ // so checkSparkAnswerAndOperator also pins that the text scan did
get bridged - an
+ // un-bridged Spark scan would read as non-native and fail here. It
is also strict-testing
+ // safe, unlike checkSparkAnswer on a fully-native plan.
+ val (_, cometPlan) = checkSparkAnswerAndOperator(sql(query))
+ // With AQE off the whole build subtree is visible in one tree, so
the bridge on the deep
+ // text scan is collectable. (Under AQE it lives inside a
materialized shuffle input stage
+ // and may not be reachable from the final plan, so only assert it
when AQE is off.)
+ if (!aqe) {
Review Comment:
The conversion can be found here under AQE. `collect` descends into query
stages, and the final plan has `CometSparkRowToColumnar` inside the build-side
shuffle stage. Could we drop the `if (!aqe)`? Under AQE,
`checkSparkAnswerAndOperator` only checks the initial plan, so right now
nothing checks the final AQE plan.
##########
spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala:
##########
@@ -301,6 +301,366 @@ class CometExecSuite extends CometTestBase {
}
}
+ test("broadcast build side with unsupported text source goes native
(#6008)") {
+ // A tiny lookup table read from a Text file (which Comet cannot scan
natively) sits on the
+ // build side of a broadcast join over a large native probe. Without the
auto-bridge the whole
+ // join stays on Spark. With
spark.comet.sparkToColumnar.broadcastBuildSide.enabled (default on)
+ // the text leaf is bridged to Arrow so the broadcast and the join run
natively.
+ assert(
+
CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.defaultValue.contains(true))
+ withTempDir { dir =>
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ spark
+ .range(0, 5)
+ .selectExpr("CAST(id AS STRING) AS value")
+ .write
+ .text(allowPath)
+ }
+
+ // Cover both the AQE and non-AQE planning paths - the tag pre-pass runs
per stage under AQE.
+ for (aqe <- Seq(false, true)) {
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe.toString) {
+ spark.read.parquet(probePath).createOrReplaceTempView("probe_6008")
+ spark.read.text(allowPath).createOrReplaceTempView("allow_6008")
+ val query =
+ "SELECT /*+ BROADCAST(a) */ p.v FROM probe_6008 p JOIN allow_6008
a ON p.key = a.value"
+
+ // Default: the text build side is bridged to Arrow and the join
goes native. The plan is
+ // fully native (the bridge counts as a native leaf), so use
checkSparkAnswerAndOperator -
+ // checkSparkAnswer would fail under COMET_STRICT_TESTING on a
fully-native plan.
+ val (_, cometPlan) = checkSparkAnswerAndOperator(sql(query))
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s
}.nonEmpty,
+ s"aqe=$aqe: expected a CometSparkToColumnarExec on the text build
side:\n" +
+ cometPlan.treeString)
+ assert(
+ collect(cometPlan) { case b: CometBroadcastExchangeExec => b
}.nonEmpty,
+ s"aqe=$aqe: expected a
CometBroadcastExchangeExec:\n${cometPlan.treeString}")
+ assert(
+ collect(cometPlan) { case j: CometBroadcastHashJoinExec => j
}.nonEmpty,
+ s"aqe=$aqe: expected a
CometBroadcastHashJoinExec:\n${cometPlan.treeString}")
+
+ // Off: the unsupported text leaf blocks native execution of the
broadcast and the join.
+ withSQLConf(
+ CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.key ->
"false") {
+ val (_, cometPlanOff) = checkSparkAnswer(sql(query))
+ assert(
+ collect(cometPlanOff) { case s: CometSparkToColumnarExec => s
}.isEmpty,
+ s"aqe=$aqe: expected no CometSparkToColumnarExec when
disabled:\n" +
+ cometPlanOff.treeString)
+ assert(
+ collect(cometPlanOff) { case b: CometBroadcastExchangeExec => b
}.isEmpty,
+ s"aqe=$aqe: expected no CometBroadcastExchangeExec when
disabled:\n" +
+ cometPlanOff.treeString)
+ assert(
+ collect(cometPlanOff) { case j: CometBroadcastHashJoinExec => j
}.isEmpty,
+ s"aqe=$aqe: expected no CometBroadcastHashJoinExec when
disabled:\n" +
+ cometPlanOff.treeString)
+ }
+ }
+ }
+ }
+ }
+
+ test("broadcast build side bridge is gated on broadcast-exchange conversion
(#6008)") {
+ // The bridge is pointless when Comet broadcast exchange conversion is
off: the broadcast (and
+ // so the join) cannot go native. With
COMET_EXEC_BROADCAST_EXCHANGE_ENABLED=false the text
+ // build side must not be bridged even though broadcastBuildSide is on.
+ withTempDir { dir =>
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ spark.range(0, 5).selectExpr("CAST(id AS STRING) AS
value").write.text(allowPath)
+ }
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.key ->
"true",
+ CometConf.COMET_EXEC_BROADCAST_EXCHANGE_ENABLED.key -> "false") {
+
spark.read.parquet(probePath).createOrReplaceTempView("probe_gate_6008")
+ spark.read.text(allowPath).createOrReplaceTempView("allow_gate_6008")
+ val (_, cometPlan) = checkSparkAnswer(sql("""
+ |SELECT /*+ BROADCAST(a) */ p.v
+ |FROM probe_gate_6008 p JOIN allow_gate_6008 a ON p.key = a.value
+ |""".stripMargin))
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s }.isEmpty,
+ "with broadcast-exchange conversion off the build side must not be
bridged:\n" +
+ cometPlan.treeString)
+ }
+ }
+ }
+
+ test("broadcast build side scan below an aggregate is bridged, shuffle stage
is not (#6008)") {
+ // A file scan can sit BELOW build-side operators. Broadcasting an
aggregated Text table puts a
+ // Text FileSourceScanExec under a partial aggregate + shuffle on the
build side. The tag
+ // pre-pass walks the whole broadcast subtree, so the deep text scan must
be bridged to Arrow,
+ // yet no query stage (e.g. the build-side shuffle under AQE) may be
wrapped in the bridge.
+ withTempDir { dir =>
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark.range(0, 20).selectExpr("CAST(id % 5 AS STRING) AS
value").write.text(allowPath)
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id % 5 AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ }
+ spark.read.text(allowPath).createOrReplaceTempView("allow_agg_6008")
+ spark.read.parquet(probePath).createOrReplaceTempView("probe_agg_6008")
+ val query = """
+ |SELECT /*+ BROADCAST(d) */ p.v, d.n
+ |FROM probe_agg_6008 p
+ |JOIN (SELECT value AS k, COUNT(*) AS n FROM allow_agg_6008 GROUP BY
value) d
+ | ON p.key = d.k
+ |""".stripMargin
+ for (aqe <- Seq(false, true)) {
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe.toString,
+ SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> (10 * 1024 *
1024).toString) {
+ // The plan is fully native (the bridged text scan under the
aggregate is a native leaf),
+ // so checkSparkAnswerAndOperator also pins that the text scan did
get bridged - an
+ // un-bridged Spark scan would read as non-native and fail here. It
is also strict-testing
+ // safe, unlike checkSparkAnswer on a fully-native plan.
+ val (_, cometPlan) = checkSparkAnswerAndOperator(sql(query))
+ // With AQE off the whole build subtree is visible in one tree, so
the bridge on the deep
+ // text scan is collectable. (Under AQE it lives inside a
materialized shuffle input stage
+ // and may not be reachable from the final plan, so only assert it
when AQE is off.)
+ if (!aqe) {
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s
}.nonEmpty,
+ s"expected the build-side text scan to be
bridged:\n${cometPlan.treeString}")
+ }
+ // A query stage (the build-side shuffle) must never be wrapped in
the bridge.
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s.child
}.foreach { child =>
+ assert(
+ !child.isInstanceOf[QueryStageExec],
+ s"aqe=$aqe: a QueryStageExec must not be bridged to
Arrow:\n${cometPlan.treeString}")
+ }
+ }
+ }
+ }
+ }
+
+ test("per-format opt-out still wins on a broadcast build side (#6008)") {
+ // A CSV build side stays a Spark FileSourceScanExec and is tagged, but
the CSV arm of
+ // shouldApplySparkToColumnar honors COMET_CONVERT_FROM_CSV_ENABLED and
never calls the
+ // broadcast-build-side bypass. So with convert-from-csv off, the CSV
build side must NOT be
+ // bridged and the join must stay on Spark, even with broadcastBuildSide
enabled.
+ withTempDir { dir =>
+ val allowPath = s"${dir.getAbsolutePath}/allow.csv"
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark.range(0, 5).selectExpr("CAST(id AS STRING) AS
value").write.csv(allowPath)
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ }
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ CometConf.COMET_CONVERT_FROM_CSV_ENABLED.key -> "false",
+ CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.key ->
"true") {
+ spark.read.csv(allowPath).createOrReplaceTempView("allow_csv_6008")
+ spark.read.parquet(probePath).createOrReplaceTempView("probe_csv_6008")
+ val (_, cometPlan) = checkSparkAnswer(sql("""
+ |SELECT /*+ BROADCAST(a) */ p.v
+ |FROM probe_csv_6008 p JOIN allow_csv_6008 a ON p.key = a._c0
+ |""".stripMargin))
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s }.isEmpty,
+ "a CSV build side with convert-from-csv disabled must not be
bridged:\n" +
+ cometPlan.treeString)
+ }
+ }
+ }
+
+ test("broadcast build side with an unsupported DSv2 scan goes native
(#6008)") {
+ // Exercises the BatchScanExec arm: an unsupported-format scan (Text read
as DSv2) on the build
+ // side is bridged so the join goes native. Keep parquet on the V1 list so
the probe stays a
+ // native CometScan; drop text from the list so it is read as a DSv2
BatchScan/TextScan.
+ withTempDir { dir =>
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark.range(0, 5).selectExpr("CAST(id AS STRING) AS
value").write.text(allowPath)
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ }
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ SQLConf.USE_V1_SOURCE_LIST.key -> "parquet") {
+ spark.read.text(allowPath).createOrReplaceTempView("allow_v2_6008")
+ spark.read.parquet(probePath).createOrReplaceTempView("probe_v2_6008")
+ // Fully native (bridged DSv2 text build side + native parquet probe +
native join), so use
+ // checkSparkAnswerAndOperator - strict-testing safe, unlike
checkSparkAnswer here.
+ val (_, cometPlan) = checkSparkAnswerAndOperator(sql("""
+ |SELECT /*+ BROADCAST(a) */ p.v
+ |FROM probe_v2_6008 p JOIN allow_v2_6008 a ON p.key = a.value
+ |""".stripMargin))
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s
}.nonEmpty,
+ s"expected the DSv2 text build side to be
bridged:\n${cometPlan.treeString}")
+ assert(
+ collect(cometPlan) { case j: CometBroadcastHashJoinExec => j
}.nonEmpty,
+ s"expected the join to go native over the DSv2 build
side:\n${cometPlan.treeString}")
+ }
+ }
+ }
+
+ test("broadcast join with an unsupported probe stays on Spark (#6008)") {
+ // Both sides are unsupported Text sources. The probe is not natively
scannable, so the join
+ // cannot go native - and because the probe is not native, the build side
must NOT be bridged
+ // either (bridging it would be useless and would break Spark's
broadcast/DPP reuse). This pins
+ // the probe-native gate that fixed the DynamicPartitionPruningV2
regression.
+ withTempDir { dir =>
+ val probePath = s"${dir.getAbsolutePath}/probe.txt"
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark.range(0, 100).selectExpr("CAST(id AS STRING) AS
value").write.text(probePath)
+ spark.range(0, 5).selectExpr("CAST(id AS STRING) AS
value").write.text(allowPath)
+ }
+ for (aqe <- Seq(false, true)) {
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe.toString) {
+ spark.read.text(probePath).createOrReplaceTempView("probe_txt_6008")
+ spark.read.text(allowPath).createOrReplaceTempView("allow_txt_6008")
+ val (_, cometPlan) = checkSparkAnswer(sql("""
+ |SELECT /*+ BROADCAST(a) */ p.value
+ |FROM probe_txt_6008 p JOIN allow_txt_6008 a ON p.value = a.value
+ |""".stripMargin))
+ // The non-native probe must leave the build side un-bridged (the
probe-native gate).
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s
}.isEmpty,
+ s"aqe=$aqe: a non-native probe must leave the build side
un-bridged:\n" +
+ cometPlan.treeString)
+ assert(
+ collect(cometPlan) { case j: CometBroadcastHashJoinExec => j
}.isEmpty,
+ s"aqe=$aqe: the join must stay on Spark when the probe is not
native:\n" +
+ cometPlan.treeString)
+ assert(
+ collect(cometPlan) { case b: CometBroadcastExchangeExec => b
}.isEmpty,
+ s"aqe=$aqe: the broadcast must stay on Spark when the join does
not convert:\n" +
+ cometPlan.treeString)
+ }
+ }
+ }
+ }
+
+ test("broadcast build side on the left (BuildLeft) with a text source goes
native (#6008)") {
+ // All other tests broadcast the right table (BuildRight). Put the small
text table first and
+ // hint it so it becomes the left build side, exercising the BuildLeft arm
of tagBuildIfProbeNative.
+ withTempDir { dir =>
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id AS STRING) AS key", "id AS v")
+ .write
+ .parquet(probePath)
+ spark.range(0, 5).selectExpr("CAST(id AS STRING) AS
value").write.text(allowPath)
+ }
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
+ spark.read.text(allowPath).createOrReplaceTempView("allow_bl_6008")
+ spark.read.parquet(probePath).createOrReplaceTempView("probe_bl_6008")
+ val (_, cometPlan) = checkSparkAnswerAndOperator(sql("""
+ |SELECT /*+ BROADCAST(a) */ p.v
+ |FROM allow_bl_6008 a JOIN probe_bl_6008 p ON a.value = p.key
+ |""".stripMargin))
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s
}.nonEmpty,
+ s"expected the text build side to be
bridged:\n${cometPlan.treeString}")
+ assert(
+ collect(cometPlan) { case j: CometBroadcastHashJoinExec => j
}.nonEmpty,
+ s"expected a CometBroadcastHashJoinExec:\n${cometPlan.treeString}")
+ }
+ }
+ }
+
+ test("broadcast nested loop join build side with a text source goes native
(#6008)") {
+ // Exercises the BroadcastNestedLoopJoinExec arm of
tagBroadcastBuildSideLeaves: an inequality
+ // join forces a BNLJ; the text build side is bridged so the BNLJ runs
natively.
+ withTempDir { dir =>
+ val probePath = s"${dir.getAbsolutePath}/probe.parquet"
+ val allowPath = s"${dir.getAbsolutePath}/allow.txt"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark.range(0, 100).selectExpr("id AS v", "CAST(id AS INT) AS
k").write.parquet(probePath)
+ spark.range(0, 5).selectExpr("CAST(id AS STRING) AS
value").write.text(allowPath)
+ }
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
+ spark.read.text(allowPath).createOrReplaceTempView("allow_bnlj_6008")
+
spark.read.parquet(probePath).createOrReplaceTempView("probe_bnlj_6008")
+ val (_, cometPlan) = checkSparkAnswer(sql("""
+ |SELECT /*+ BROADCAST(a) */ p.v
+ |FROM probe_bnlj_6008 p JOIN allow_bnlj_6008 a ON p.k >
CAST(a.value AS INT)
+ |""".stripMargin))
+ assert(
+ collect(cometPlan) { case s: CometSparkToColumnarExec => s
}.nonEmpty,
+ s"expected the text build side to be
bridged:\n${cometPlan.treeString}")
+ assert(
+ collect(cometPlan) { case j: CometBroadcastNestedLoopJoinExec => j
}.nonEmpty,
+ s"expected a
CometBroadcastNestedLoopJoinExec:\n${cometPlan.treeString}")
+ }
+ }
+ }
+
+ test("AQE DPP with a text dimension on the broadcast build side runs
natively (#6008)") {
+ // A Text dimension on a broadcast build side joined to a PARTITIONED
Parquet fact triggers DPP.
+ // The build-side bridge makes the join a CometBroadcastHashJoinExec; the
DPP subquery's own
+ // build (the text scan) is not bridged, so
CometPlanAdaptiveDynamicPruningFilters must NOT wrap
+ // it in a CometBroadcastExchangeExec (that crashed with "Comet execution
only takes Arrow
+ // Arrays"). It falls back to a Spark DPP broadcast while the join stays
native. See issue #6008.
+ assume(isSpark35Plus, "Native AQE DPP requires Spark 3.5+")
+ withSQLConf(
+ SQLConf.USE_V1_SOURCE_LIST.key -> "parquet",
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "true",
+ SQLConf.DYNAMIC_PARTITION_PRUNING_ENABLED.key -> "true",
+ SQLConf.DYNAMIC_PARTITION_PRUNING_REUSE_BROADCAST_ONLY.key -> "true") {
+ withTempDir { dir =>
+ val factPath = s"${dir.getAbsolutePath}/fact.parquet"
+ val dimPath = s"${dir.getAbsolutePath}/dim.txt"
+ withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
+ spark
+ .range(0, 100)
+ .selectExpr("CAST(id % 10 AS STRING) AS fk", "id AS v")
+ .write
+ .partitionBy("fk")
+ .parquet(factPath)
+ spark.range(0, 10).selectExpr("CAST(id AS STRING) AS
value").write.text(dimPath)
+ }
+ spark.read.parquet(factPath).createOrReplaceTempView("dpp_fact_txt")
+ spark.read.text(dimPath).createOrReplaceTempView("dpp_dim_txt")
+ val df = sql("""
+ |SELECT /*+ BROADCAST(d) */ f.v
+ |FROM dpp_fact_txt f JOIN dpp_dim_txt d ON f.fk = d.value
+ |WHERE d.value < '5'
Review Comment:
Spark infers `d.value < '5'` onto `f.fk`, so the fact is pruned statically
and DPP adds nothing here. This test would still pass if DPP were replaced by
`true`. Could the dimension filter be something Spark can't push to the fact,
like `CAST(d.value AS INT) % 2 = 0`? And could the test assert that the fact
scan still has a live DPP subquery?
##########
spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:
##########
@@ -1518,6 +1535,126 @@ case class CometExecRule(session: SparkSession,
queryStagePrep: Boolean = false)
}
}
+ /**
+ * An unsupported file-source scan is bridged to Arrow when the general
sparkToColumnar opt-in
+ * covers it, or when it feeds a broadcast join's build side (see issue
#6008). This is only
+ * reached from the unsupported-format arms, so the per-format
`COMET_CONVERT_FROM_*` opt-outs
+ * for CSV/JSON/Parquet still take precedence.
+ */
+ private def shouldBridgeUnsupportedScan(conf: SQLConf, op: SparkPlan):
Boolean =
+ isSparkToArrowEnabled(conf, op) || isBroadcastBuildSideLeaf(op)
+
+ /**
+ * True when the build-side auto-bridge should run: its own config is on and
Comet broadcast
+ * exchange conversion is enabled. Without the latter the
`BroadcastExchangeExec` cannot become
+ * a `CometBroadcastExchangeExec`, so bridging the build side would only add
a row to Arrow copy
+ * for no native join. This does not also check the join serdes
+ * (COMET_EXEC_BROADCAST_HASH_JOIN_ENABLED /
COMET_EXEC_BROADCAST_NESTED_LOOP_JOIN_ENABLED): if
+ * the exchange converts but the join does not, the broadcast arm returns
the original plan and
+ * EliminateRedundantTransitions removes the now-unused bridge, so a
leftover copy never reaches
+ * runtime. See issue #6008.
+ */
+ private def broadcastBuildSideBridgeEnabled: Boolean =
+ CometConf.COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED.get(conf) &&
+ CometConf.COMET_EXEC_BROADCAST_EXCHANGE_ENABLED.get(conf)
+
+ /**
+ * True when `op` is a scan on a broadcast join's build side and build-side
bridging is enabled.
+ * `tagBroadcastBuildSideLeaves` only sets the tag when bridging is enabled,
but the tag lives
+ * on the physical node and can survive onto a reused instance, so re-check
here: a stale tag
+ * must not bridge a scan once the feature is turned off. A tagged scan is
bridged to Arrow even
+ * when the general sparkToColumnar path is off, letting the build branch
and the join run
+ * natively. See issue #6008.
+ */
+ private def isBroadcastBuildSideLeaf(op: SparkPlan): Boolean =
+ broadcastBuildSideBridgeEnabled &&
+ op.getTagValue(CometExecRule.BROADCAST_BUILD_SIDE_TAG).isDefined
+
+ /**
+ * Tag the file-source scan leaves on a broadcast join's build side so
+ * shouldApplySparkToColumnar can bridge an unsupported build-side scan
(e.g. a Text scan) to
+ * Arrow, letting the whole join run natively. The build side is usually
small (auto-broadcasts
+ * are capped by the broadcast threshold), so the row to Arrow copy is
usually cheap - but that
+ * is a heuristic, not a bound: an explicit BROADCAST hint can force a
larger build side, and a
+ * selective filter above the scan means the bridge copies the pre-filter
scan output. See issue
+ * #6008.
+ *
+ * The bridge is applied ONLY when the probe side is already natively
scannable, i.e. the join
+ * can actually become a fully native CometBroadcastHashJoinExec. Bridging a
build side under a
+ * join that stays on Spark is pure overhead and, worse, it rewrites the
build broadcast's
+ * subtree to Comet and breaks Spark's DPP broadcast reuse (the DPP subquery
keeps an unbridged
+ * copy of the same scan, so the two exchanges no longer share a canonical
form). CometScanRule
+ * has already run by this point, so a native probe scan is a Comet scan and
an unsupported
+ * probe scan is still a plain Spark scan - see `hasOnlyNativeScans`.
+ *
+ * Only `FileSourceScanExec` / `BatchScanExec` are tagged - the scan types
the per-format arms
+ * of shouldApplySparkToColumnar handle. A natively scannable file (e.g.
Parquet) is normally
+ * already a `CometScanExec` here; if native scan conversion did not take,
its
+ * `FileSourceScanExec` is still tagged, but the tag is harmless because the
per-format arm
+ * decides that scan before isBroadcastBuildSideLeaf is ever consulted. DPP
subquery broadcasts
+ * live in expressions, not `children`, so `foreach` never reaches them -
that scope is
+ * intentional.
+ */
+ private def tagBroadcastBuildSideLeaves(plan: SparkPlan): Unit = {
+ if (!broadcastBuildSideBridgeEnabled) {
+ return
+ }
+ plan.foreach {
+ case j: BroadcastHashJoinExec => tagBuildIfProbeNative(j.buildSide,
j.left, j.right)
+ case j: BroadcastNestedLoopJoinExec =>
tagBuildIfProbeNative(j.buildSide, j.left, j.right)
+ case _ =>
+ }
+ }
+
+ /** Tag the build side's file-source scans, but only if the probe side is
natively scannable. */
+ private def tagBuildIfProbeNative(
+ buildSide: BuildSide,
+ left: SparkPlan,
+ right: SparkPlan): Unit = {
+ val (buildPlan, probePlan) = buildSide match {
+ case BuildLeft => (left, right)
+ case BuildRight => (right, left)
+ }
+ if (hasOnlyNativeScans(probePlan)) {
+ tagBuildSideScans(buildPlan)
+ }
+ }
+
+ /**
+ * Tag the broadcast dimension's own file-source scans, descending through
its filters, projects
+ * and aggregations (incl. the aggregate's shuffle) but stopping at a nested
join. A nested
+ * join's inputs are not the bounded broadcast dimension - its streamed side
can be far larger
+ * than the broadcast output - so bridging them would copy a large input to
Arrow for no reason.
+ * That nested join is tagged by its own pre-pass visit based on its own
probe. See issue #6008.
+ */
+ private def tagBuildSideScans(plan: SparkPlan): Unit = plan match {
+ case scan @ (_: FileSourceScanExec | _: BatchScanExec) =>
+ scan.setTagValue(CometExecRule.BROADCAST_BUILD_SIDE_TAG, ())
+ case _: BroadcastHashJoinExec | _: BroadcastNestedLoopJoinExec | _:
ShuffledHashJoinExec |
+ _: SortMergeJoinExec =>
+ case other => other.children.foreach(tagBuildSideScans)
+ }
+
+ /**
+ * True when `plan` has at least one scan and none of its scan leaves is an
unsupported plain
+ * Spark scan - i.e. every scan is already a Comet scan (CometScanRule ran
before this rule).
+ * Used to decide whether a broadcast join's probe side can feed a native
join. A probe that
+ * itself contains an unsupported scan (e.g. a DSv2
`InMemoryTableWithV2Filter` fact) keeps the
+ * join on Spark, so bridging its build side would only break DPP reuse.
+ */
+ private def hasOnlyNativeScans(plan: SparkPlan): Boolean = {
Review Comment:
With two text dimensions on one fact, only the inner join goes native. When
the pre-pass checks the outer join, its probe side still contains the inner
join's text scan, which isn't converted yet, so `hasOnlyNativeScans` returns
false. Is that intended? If it's out of scope here, could you open an issue for
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]