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]

Reply via email to