Repository navigation
feat: run broadcast joins natively when an unsupported scan is on the build side - #6786
parthchandra wants to merge 2 commits into
Conversation
|
This seems like a great idea. Also, |
andygrove
left a comment
There was a problem hiding this comment.
Thanks for this. Getting the join over the large probe side to run natively when a small lookup table is in a format Comet can't scan is a nice win. The ORC, Iceberg and text probes I ran all matched Spark. I did find three cases where the change makes an existing query slower than it is today, so I'm requesting changes for those. Details are inline. The new tests pass for me on 3.4, 3.5 and 4.2 as well.
Since this changes the planner, I've added run-spark-4.1-tests and run-all-spark-profiles. The local sql_core-1 run skips org.apache.spark.sql.execution.datasources.* and org.apache.spark.sql.connector.*, which run in sql_core-4. Those are the suites that cover ORC, text and V2 scans, and Hive ORC tables are covered in sql_hive. The profiles label runs the new tests on 3.4, where the DPP path is different.
docs/source/user-guide/latest/datasources.md still says these scans only convert with spark.comet.sparkToColumnar.enabled and the operator list. Could you add a paragraph there about the new default and how to turn it off?
The earlier [scans] failure was the listener race in the Iceberg write report test, which #6809 just fixed on main. The re-run passed.
| case BuildLeft => (left, right) | ||
| case BuildRight => (right, left) | ||
| } | ||
| if (hasOnlyNativeScans(probePlan)) { |
There was a problem hiding this 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.
There was a problem hiding this comment.
Good catch. tagBuildIfProbeNative now skips a build side whose BroadcastExchangeExec has SKIP_COMET_BROADCAST_TAG, so on 3.4 the text scan isn't bridged and DPP stays on.
There was a problem hiding this comment.
Thanks, this works now. On 3.4.3 with AQE the fact scan reads 400 rows again and the broadcast is reused. The new DPP test is skipped below 3.5, though, so nothing would catch it if this came back. Could you add a 3.4 test with a native Iceberg fact partitioned on the join key and a text dimension, filtered with something Spark can't infer onto the fact, that asserts the fact's DPP filter isn't true? The last test in the probe suite I posted in the ORC thread does this.
| * 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 |
There was a problem hiding this 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.
There was a problem hiding this comment.
Agreed. When the join doesn't convert, the rule now reverts the build side to its original Spark plan (revertBridgedBroadcasts). That handles the FULL/LEFT OUTER and broadcastHashJoin.enabled=false cases, and the 3.4 case too, so the dimension isn't broadcast twice.
There was a problem hiding this comment.
I don't think revertBridgedBroadcasts ever fires. During the transform the bridge is CometScanWrapper(nativeOp, CometSparkToColumnarExec(scan)), and CometScanWrapper is a leaf, so b.exists(_.isInstanceOf[CometSparkToColumnarExec]) never sees it. On 3738210b84 the FULL OUTER and build-left LEFT OUTER nested loop joins still end up with BroadcastExchange <- CometColumnarToRow <- CometFilter <- CometSparkRowToColumnar <- FileScan text, with AQE on and off. The same happens when the hash join declines because an operator on the probe side stays on Spark. I used an explode with spark.comet.exec.explode.enabled=false to get one. With DPP on top, that query scans the text dimension twice and builds two broadcasts. With the config off it scans it once and reuses the broadcast.
I tried making the check look inside CometScanWrapper, and that isn't enough on its own. With AQE off, a filtered dimension goes back to plain Spark and DPP reuses the broadcast again. With AQE on, the stage gets prepared again without the join above it, and the scan still has BROADCAST_BUILD_SIDE_TAG, so the bridge is back in the final plan. revertCometToSpark also rebuilds both a CometSinkPlaceHolder and the Comet node inside it from the same originalPlan. With an aggregate on the build side, that gives a Spark Exchange over a CometExchange over a Spark partial aggregate, which fails at runtime with ClassCastException: OffHeapColumnVector cannot be cast to CometVector. With a UNION ALL on the build side, planning fails with Incorrect number of children. It also skips the sparkFallback overrides in CometNativeScanExec, CometIcebergNativeScanExec and CometLocalTopKExec.
Could the rule remove the tag from the scans under a join that stays on Spark and convert that build side again, instead of rebuilding the Spark nodes? That should give exactly the plan we get today, and with the tag gone the AQE stage pass won't bridge the scan again. Could you also add tests for a FULL OUTER nested loop join in both AQE modes, with a filter, an aggregate and a union on the build side, that assert there's no CometSparkToColumnarExec under a Spark BroadcastExchangeExec? A DPP query where the join declines could assert the ReusedExchangeExec. The docstring here on broadcastBuildSideBridgeEnabled also still says EliminateRedundantTransitions removes the unused bridge and that the join configs aren't checked.
| } | ||
| case other => other | ||
| } | ||
| stripped.isInstanceOf[CometNativeExec] |
There was a problem hiding this 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.
There was a problem hiding this comment.
Fixed — changed it to stripped.isInstanceOf[CometPlan] && stripped.supportsColumnar, so CometUnion/CometCoalesce/CometTakeOrderedAndProject are accepted and the UNION-ALL dimension keeps its reuse. Thanks for filing #6815 for the non-AQE twin.
There was a problem hiding this comment.
Thanks, the union dimension gets its ReusedExchange back under AQE on 4.1. Could you add a test so it stays that way? An AQE DPP query whose dimension is a UNION ALL of two filtered Parquet tables, asserting a ReusedExchangeExec under the DPP subquery, fails with the old CometNativeExec check. The docstring above isNativeBuildSide also still says it checks for a CometNativeExec.
| } | ||
| } | ||
|
|
||
| test("broadcast build side with unsupported text source goes native (#6008)") { |
There was a problem hiding this 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.
There was a problem hiding this comment.
Yes please, I'll take the probe. I can add an ORC build-side test (nulls across types, struct, ARRAY, MAP, timestamp/ntz, LA timezone) and an Iceberg-declined one from it.
There was a problem hiding this comment.
Here it is. It runs as its own suite under spark/src/test/scala/org/apache/comet/. Each query goes through checkSparkAnswer, and the printed line shows whether the bridge and the native join were planned. All three tests pass on 3738210b84 with 3.4, and the first two also pass on 4.1. For the PR, asserting the bridge and the CometBroadcastHashJoinExec instead of printing them should be enough.
Pr6786BuildSideProbeSuite.scala
package org.apache.comet
import org.apache.spark.sql.{CometTestBase, DataFrame}
import org.apache.spark.sql.catalyst.expressions.{DynamicPruningExpression, Literal}
import org.apache.spark.sql.comet._
import org.apache.spark.sql.execution._
import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanExec, QueryStageExec}
import org.apache.spark.sql.execution.exchange.ReusedExchangeExec
import org.apache.spark.sql.internal.SQLConf
// Probes from the PR #6786 review: ORC and Iceberg build sides that Comet cannot scan natively,
// bridged by spark.comet.convert.broadcastBuildSide.enabled. Each query is compared with Spark
// through checkSparkAnswer, and the summary line shows whether the bridge and the native join
// were planned.
class Pr6786BuildSideProbeSuite extends CometTestBase with CometIcebergTestBase {
/** Every node, descending through AQE, stages and subqueries, but not into reused exchanges. */
private def allNodes(p: SparkPlan): Seq[SparkPlan] = {
val kids = p match {
case a: AdaptiveSparkPlanExec => Seq(a.executedPlan)
case s: QueryStageExec => Seq(s.plan)
case _: ReusedExchangeExec => Seq.empty
case other => other.children
}
val subs = p.expressions.flatMap(_.collect { case e: ExecSubqueryExpression => e.plan })
Seq(p) ++ kids.flatMap(allNodes) ++ subs.flatMap(allNodes)
}
private def summary(df: DataFrame): String = {
df.collect()
val nodes = allNodes(df.queryExecution.executedPlan)
s"bridges=${nodes.count(_.isInstanceOf[CometSparkToColumnarExec])} " +
s"cometBHJ=${nodes.count(_.isInstanceOf[CometBroadcastHashJoinExec])} " +
s"icebergNative=${nodes.count(_.isInstanceOf[CometIcebergNativeScanExec])}"
}
private def dppSummary(df: DataFrame): String = {
df.collect()
val nodes = allNodes(df.queryExecution.executedPlan)
val dppTrue = nodes.exists(_.expressions.exists(_.exists {
case DynamicPruningExpression(Literal.TrueLiteral) => true
case _ => false
}))
val dppLive = nodes.exists(_.expressions.exists(_.exists {
case DynamicPruningExpression(e) if e != Literal.TrueLiteral => true
case _ => false
}))
summary(df) + s" dppLive=$dppLive dppTrue=$dppTrue " +
s"reusedX=${nodes.count(_.isInstanceOf[ReusedExchangeExec])}"
}
test("ORC build side, all types, non-UTC session (#6008)") {
withTempDir { dir =>
val fact = s"${dir.getAbsolutePath}/fact.parquet"
val dim = s"${dir.getAbsolutePath}/dim.orc"
withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
spark.range(0, 2000).selectExpr("CAST(id % 60 AS INT) AS k", "id AS v").write.parquet(fact)
spark
.range(0, 60, 1, 3)
.selectExpr(
"CAST(id AS INT) AS k",
"IF(id % 7 = 0, NULL, CAST(id AS BYTE)) AS b",
"IF(id % 7 = 1, NULL, CAST(id AS SHORT)) AS s",
"IF(id % 7 = 2, NULL, id * 1000000000L) AS l",
"IF(id % 7 = 3, NULL, id % 2 = 0) AS bo",
"IF(id % 7 = 4, NULL, CAST(id AS FLOAT) / 3) AS f",
"CASE WHEN id % 7 = 5 THEN NULL WHEN id = 6 THEN CAST('NaN' AS DOUBLE) " +
"WHEN id = 8 THEN -0.0D ELSE id / 3.0D END AS d",
"IF(id % 7 = 6, NULL, CAST(id * 1.25 AS DECIMAL(10,2))) AS dec1",
"IF(id % 5 = 0, NULL, CAST(id AS DECIMAL(38,10)) / 7) AS dec2",
"IF(id % 5 = 1, NULL, CONCAT('s', id, 'é')) AS str",
"IF(id % 5 = 2, NULL, CAST(CONCAT('b', id) AS BINARY)) AS bin",
"IF(id % 5 = 3, NULL, DATE_ADD(DATE'1960-01-01', CAST(id * 300 AS INT))) AS dt",
"IF(id % 5 = 4, NULL, TIMESTAMP_MICROS(id * 8640000000000L - 900000000000000L)) AS ts",
"IF(id % 9 = 0, NULL, " +
"CAST(TIMESTAMP_MICROS(id * 8640000000000L - 900000000000000L) AS TIMESTAMP_NTZ)) " +
"AS ntz",
"IF(id % 6 = 0, NULL, named_struct('a', CAST(id AS INT), 'b', " +
"IF(id % 4 = 0, NULL, CONCAT('x', id)))) AS st",
"CASE WHEN id % 6 = 1 THEN NULL WHEN id % 6 = 2 THEN CAST(array() AS ARRAY<STRING>) " +
"ELSE array(CONCAT('e', id), NULL, 'z') END AS arr",
"CASE WHEN id % 6 = 3 THEN NULL WHEN id % 6 = 4 THEN " +
"CAST(map() AS MAP<STRING,STRING>) " +
"ELSE map(CONCAT('k', id), IF(id % 2 = 0, NULL, 'v')) END AS m")
.write
.orc(dim)
}
val query =
"""SELECT /*+ BROADCAST(d) */ f.v, d.*, hour(d.ts) AS h, CAST(d.ts AS STRING) AS tss,
| CAST(d.ntz AS STRING) AS ntzs, date_trunc('DAY', d.ts) AS tday
|FROM fact_o f JOIN dim_o d ON f.k = d.k""".stripMargin
for {
aqe <- Seq("false", "true")
v1List <- Seq("avro,csv,json,kafka,orc,parquet,text", "parquet")
(vec, nestedVec) <- Seq(("true", "true"), ("true", "false"), ("false", "false"))
} {
withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe,
SQLConf.USE_V1_SOURCE_LIST.key -> v1List,
SQLConf.ORC_VECTORIZED_READER_ENABLED.key -> vec,
SQLConf.ORC_VECTORIZED_READER_NESTED_COLUMN_ENABLED.key -> nestedVec,
SQLConf.SESSION_LOCAL_TIMEZONE.key -> "America/Los_Angeles") {
spark.read.parquet(fact).createOrReplaceTempView("fact_o")
spark.read.orc(dim).createOrReplaceTempView("dim_o")
// scalastyle:off println
println(s"ORC aqe=$aqe v1=$v1List vec=$vec nested=$nestedVec: " +
summary(sql(query)))
// scalastyle:on println
checkSparkAnswer(sql(query))
}
}
}
}
test("Iceberg build side that Comet's native reader declines (#6008)") {
assume(icebergAvailable, "Iceberg not available in classpath")
withTempIcebergDir { warehouseDir =>
withTempDir { dir =>
val fact = s"${dir.getAbsolutePath}/fact.parquet"
withSQLConf(
"spark.sql.catalog.pc" -> "org.apache.iceberg.spark.SparkCatalog",
"spark.sql.catalog.pc.type" -> "hadoop",
"spark.sql.catalog.pc.warehouse" -> warehouseDir.getAbsolutePath) {
withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
spark
.range(0, 2000)
.selectExpr("CAST(id % 40 AS INT) AS k", "id AS v")
.write
.parquet(fact)
for ((name, props) <- Seq(
"dim_orc" -> "'write.format.default'='orc'",
"dim_pq" -> ("'format-version'='2', 'write.delete.mode'='merge-on-read', " +
"'write.update.mode'='merge-on-read'"))) {
spark.sql(s"""CREATE TABLE pc.db.$name (
| k INT, s STRING, d DECIMAL(10,2), ts TIMESTAMP,
| st STRUCT<a: INT, b: STRING>, arr ARRAY<STRING>,
| m MAP<STRING, STRING>, region STRING)
|USING iceberg PARTITIONED BY (region)
|TBLPROPERTIES ($props)""".stripMargin)
spark
.range(0, 40, 1, 2)
.selectExpr(
"CAST(id AS INT) AS k",
"IF(id % 5 = 0, NULL, CONCAT('s', id)) AS s",
"IF(id % 6 = 0, NULL, CAST(id * 1.5 AS DECIMAL(10,2))) AS d",
"IF(id % 7 = 0, NULL, " +
"TIMESTAMP_MICROS(id * 8640000000000L - 900000000000000L)) AS ts",
"IF(id % 4 = 0, NULL, named_struct('a', CAST(id AS INT), 'b', " +
"IF(id % 3 = 0, NULL, CONCAT('b', id)))) AS st",
"CASE WHEN id % 4 = 1 THEN NULL WHEN id % 4 = 2 THEN " +
"CAST(array() AS ARRAY<STRING>) ELSE array(CONCAT('e', id), NULL) END AS arr",
"IF(id % 3 = 1, NULL, map(CONCAT('k', id), IF(id % 2 = 0, NULL, 'v'))) AS m",
"IF(id % 2 = 0, 'east', 'west') AS region")
.writeTo(s"pc.db.$name")
.append()
}
spark.sql("DELETE FROM pc.db.dim_pq WHERE k % 9 = 0")
}
spark.read.parquet(fact).createOrReplaceTempView("ice_fact")
val queries = Seq(
"orc" -> ("SELECT /*+ BROADCAST(d) */ f.v, d.* FROM ice_fact f " +
"JOIN pc.db.dim_orc d ON f.k = d.k"),
"pq-mor" -> ("SELECT /*+ BROADCAST(d) */ f.v, d.* FROM ice_fact f " +
"JOIN pc.db.dim_pq d ON f.k = d.k"),
"snapshots" -> ("SELECT /*+ BROADCAST(s) */ f.v, s.operation, s.summary FROM " +
"ice_fact f JOIN pc.db.dim_pq.snapshots s ON f.k = CAST(s.snapshot_id % 2 AS INT)"))
for {
(name, query) <- queries
native <- Seq("true", "false")
aqe <- Seq("false", "true")
} {
withSQLConf(
CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> native,
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe,
SQLConf.SESSION_LOCAL_TIMEZONE.key -> "Asia/Kolkata") {
// scalastyle:off println
println(s"ICEBERG $name icebergNative=$native aqe=$aqe: ${summary(sql(query))}")
// scalastyle:on println
checkSparkAnswer(sql(query))
}
}
}
}
}
}
test("Spark 3.4: native Iceberg fact, text dimension, AQE DPP (#6008)") {
assume(icebergAvailable, "Iceberg not available in classpath")
withTempIcebergDir { warehouseDir =>
withTempDir { dir =>
val dim = s"${dir.getAbsolutePath}/dim.txt"
withSQLConf(
"spark.sql.catalog.pc" -> "org.apache.iceberg.spark.SparkCatalog",
"spark.sql.catalog.pc.type" -> "hadoop",
"spark.sql.catalog.pc.warehouse" -> warehouseDir.getAbsolutePath) {
withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
spark.sql("""CREATE TABLE pc.db.fact (v BIGINT, fk STRING)
|USING iceberg PARTITIONED BY (fk)""".stripMargin)
spark
.range(0, 1000, 1, 2)
.selectExpr("id AS v", "CAST(id % 10 AS STRING) AS fk")
.writeTo("pc.db.fact")
.append()
spark.range(0, 10).selectExpr("CAST(id AS STRING) AS value").write.text(dim)
}
spark.read.text(dim).createOrReplaceTempView("ice_dim")
val query =
"""SELECT /*+ BROADCAST(d) */ f.v FROM pc.db.fact f JOIN
|(SELECT value FROM ice_dim WHERE CAST(value AS INT) % 3 = 0) d ON f.fk = d.value""".stripMargin
for {
aqe <- Seq("true", "false")
feature <- Seq("true", "false")
} {
withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe,
SQLConf.DYNAMIC_PARTITION_PRUNING_ENABLED.key -> "true",
SQLConf.DYNAMIC_PARTITION_PRUNING_REUSE_BROADCAST_ONLY.key -> "true",
"spark.comet.convert.broadcastBuildSide.enabled" -> feature) {
// scalastyle:off println
val df = sql(query)
val dpp = dppSummary(df) // runs the query
val factRows = allNodes(df.queryExecution.executedPlan).collect {
case s: CometIcebergNativeScanExec => s.metrics("output_rows").value
}
// With DPP live the fact scan reads 400 of its 1000 rows.
println(s"DPP aqe=$aqe feature=$feature: $dpp factRows=${factRows.mkString(",")}")
// scalastyle:on println
checkSparkAnswer(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) { |
There was a problem hiding this 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.
There was a problem hiding this comment.
Done — dropped it; the bridge is asserted in both AQE modes now, so the final AQE plan is covered.
| 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' |
There was a problem hiding this 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?
There was a problem hiding this comment.
Right. Filter is now CAST(d.value AS INT) % 2 = 0 so Spark can't push it to the fact, and the test asserts the fact scan still has a live DPP SubqueryBroadcast.
| .createWithDefault(Nil) | ||
|
|
||
| val COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED: ConfigEntry[Boolean] = | ||
| conf("spark.comet.sparkToColumnar.broadcastBuildSide.enabled") |
There was a problem hiding this 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?
There was a problem hiding this comment.
Renamed to spark.comet.convert.broadcastBuildSide.enabled. Reworded the doc too: it only covers file/V2 scans, only when every probe-side scan is native, and notes the dimension is read twice under DPP.
| * 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 = { |
There was a problem hiding this 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?
There was a problem hiding this comment.
I wanted to keep this minimal so I would say this is out of scope for this PR, opened #6832 to track it. It's a missed optimization (only the inner join goes native), not wrong results.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: An unsupported build-side scan prevented native broadcast joins even when the probe used a native scan. This PR enables that execution path automatically.
- Design approach: Tag eligible build-side scans before conversion, reuse
CometSparkToColumnarExec, and guard AQE DPP broadcasts against non-Arrow inputs. - Correctness: Reviewed the conversion, schema, input-file metadata, join eligibility and DPP paths against Spark sources. Existing join restrictions remain enforced. No additional introduced P1/P2 correctness issue was found beyond the previously reported blockers.
- Compatibility analysis: Compared relevant Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Spark's DPP reuse depends on matching broadcast plans. The new scan tagging conflicts with the existing Spark 3.4 fallback, as already reported. Compared the touched paths with
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1. Automatic conversion is intentional. Lost pruning and broadcast reuse are unintended regressions. - Key design decisions: The feature defaults to enabled independently of general Spark-to-Arrow conversion. Broadcast-exchange conversion must also be enabled. CSV, JSON and Parquet conversion switches retain precedence, and tagging stops at nested joins.
- Implementation sketch:
CometExecRuletags build-side file/V2 scans when the probe contains only recognized native scans. The existing bridge converts their output.CometPlanAdaptiveDynamicPruningFilterschecks the DPP build separately before choosing its exchange type. Nine tests cover positive paths, configuration gates, aggregation, V2 scans, build orientation and DPP execution. - Performance: Conversion copies scan output before filters and aggregates, so broadcast output size does not bound conversion work. The existing review substantiates lost pruning, retained conversions after join fallback, and duplicate DPP broadcasts. No benchmark establishes a net speedup for those cases, and I make no additional performance claim.
- Design: Reusing the existing bridge keeps the execution mechanism simple. However, native scan eligibility does not guarantee native join eligibility. The existing request to restore the build subtree when the join declines addresses that mismatch.
- Abstraction & complexity: The helpers are narrowly scoped and introduce no new execution framework. The new
CometNativeExectest conflates class membership with Arrow-compatible output and excludes valid columnar operators such asCometUnionExec. This is covered by an existing blocker. - Behavioral changes worth calling out: Eligible unsupported file and V2 scans can now feed native broadcast joins by default. DPP can execute a separate Spark broadcast for an unbridged dimension. The new AQE guard also affects existing native queries when the feature switch is disabled.
- Suggested improvements: Resolve the existing requests to honor the Spark 3.4 broadcast skip tag, undo conversions when a join declines, and recognize compatible Comet columnar DPP builds. Regression assertions should verify pruning and exchange reuse. No additional improvement meeting the P1/P2 reporting bar was identified.
Reviewed the full four-file PR diff for 8932aad407544a51fbfc38167c3bb46f29a67946 against supplied base e8c7ee739d3530b309d38388ba1c41decdb8abea, using merge base f9dd86d482d136426990186c2f1d42eef6a1ca59. GitHub's changed-file list matches this scope. The PR remains non-draft. Read the review, conversation comment, eight inline comments and eight unresolved threads before forming conclusions.
Routed skills: review-comet-pr and review-comet-ffi-pr. Also checked contributor guidance on FFI ownership, memory accounting, timezones, configuration, development and CI.
Existing blockers remain substantiated by the reported reproductions and source tracing: Spark 3.4 DPP loss, retained conversion and lost reuse after join fallback, and AQE DPP rejecting valid columnar builds. These are not duplicated as new findings.
Exact-head CI: Latest check runs associated with the reviewed SHA show 60 successes and 53 skips, with none failed or running. Comet suites passed across Spark 3.4–4.2, and all eight Spark 4.1 SQL/Hive shards passed. Retrieved execution logs confirm all nine new tests passed on 4.1. On 3.4, eight passed and the explicitly unsupported AQE DPP test was canceled. These jobs tested merge commits containing this head. The default job merged the supplied base, while the later profile run merged newer main.
Validation limits: No local JVM/native tests or benchmarks were run. This checkout has no compiled artifacts or Spark dependency cache. Runtime evidence comes from CI and the existing review's reproductions, corroborated by source inspection. macOS, dedicated Iceberg suites and other versions' full Spark SQL suites were skipped. git diff --check passed.
No additional introduced P1/P2 issues found within this review beyond the existing unresolved concerns.
parthchandra
left a comment
There was a problem hiding this comment.
Thank you for the review @andygrove. Addressed your comments and also added a paragraph to datasources.md about the new default and how to turn it off
| case BuildLeft => (left, right) | ||
| case BuildRight => (right, left) | ||
| } | ||
| if (hasOnlyNativeScans(probePlan)) { |
There was a problem hiding this comment.
Good catch. tagBuildIfProbeNative now skips a build side whose BroadcastExchangeExec has SKIP_COMET_BROADCAST_TAG, so on 3.4 the text scan isn't bridged and DPP stays on.
| * 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 |
There was a problem hiding this comment.
Agreed. When the join doesn't convert, the rule now reverts the build side to its original Spark plan (revertBridgedBroadcasts). That handles the FULL/LEFT OUTER and broadcastHashJoin.enabled=false cases, and the 3.4 case too, so the dimension isn't broadcast twice.
| } | ||
| case other => other | ||
| } | ||
| stripped.isInstanceOf[CometNativeExec] |
There was a problem hiding this comment.
Fixed — changed it to stripped.isInstanceOf[CometPlan] && stripped.supportsColumnar, so CometUnion/CometCoalesce/CometTakeOrderedAndProject are accepted and the UNION-ALL dimension keeps its reuse. Thanks for filing #6815 for the non-AQE twin.
| } | ||
| } | ||
|
|
||
| test("broadcast build side with unsupported text source goes native (#6008)") { |
There was a problem hiding this comment.
Yes please, I'll take the probe. I can add an ORC build-side test (nulls across types, struct, ARRAY, MAP, timestamp/ntz, LA timezone) and an Iceberg-declined one from it.
| // 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) { |
There was a problem hiding this comment.
Done — dropped it; the bridge is asserted in both AQE modes now, so the final AQE plan is covered.
| 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' |
There was a problem hiding this comment.
Right. Filter is now CAST(d.value AS INT) % 2 = 0 so Spark can't push it to the fact, and the test asserts the fact scan still has a live DPP SubqueryBroadcast.
| .createWithDefault(Nil) | ||
|
|
||
| val COMET_SPARK_TO_ARROW_BROADCAST_BUILD_SIDE_ENABLED: ConfigEntry[Boolean] = | ||
| conf("spark.comet.sparkToColumnar.broadcastBuildSide.enabled") |
There was a problem hiding this comment.
Renamed to spark.comet.convert.broadcastBuildSide.enabled. Reworded the doc too: it only covers file/V2 scans, only when every probe-side scan is native, and notes the dimension is read twice under DPP.
| * 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 = { |
There was a problem hiding this comment.
I wanted to keep this minimal so I would say this is out of scope for this PR, opened #6832 to track it. It's a missed optimization (only the inner join goes native), not wrong results.
8932aad to
09fb41f
Compare
… build side (apache#6008) 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.
09fb41f to
3738210
Compare
andygrove
left a comment
There was a problem hiding this comment.
Thanks for the quick turnaround on these. I re-ran my probes against 3738210b84. The 3.4 fix works: with AQE the Iceberg fact reads 400 rows instead of 1000, and the broadcast is reused. The union dimension gets its ReusedExchange back under AQE, and turning off spark.comet.exec.broadcastHashJoin.enabled no longer leaves a bridge behind. The ORC and Iceberg probes still match Spark on 3.4 and 4.1.
The revert for joins that stay on Spark doesn't take effect yet, though. A FULL OUTER nested loop join still keeps the bridge under a Spark broadcast, and with DPP the dimension is still read twice. That's the slowdown from my first review, so I'm keeping the request for changes. Details are in the thread.
None of the three fixes has a test yet, so I've asked for one in each thread. CI is green, but none of the suites reaches the revert path.
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 aCometBroadcastExchangeExec,so the
BroadcastHashJoincan't go native — even when the large probe side is a fully nativescan+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 aCometScanWrapperthat is a
CometNativeExec), but it is off by default and itssupportedOperatorListexcludes filescans, 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(defaulttrue,CometConf.scala). Independent of the generalspark.comet.sparkToColumnar.enabledopt-in.Auto-bridge on broadcast build sides (
CometExecRule.scala): a pre-pass(
tagBroadcastBuildSideLeaves) tags the build-sideFileSourceScanExec/BatchScanExecof aBroadcastHashJoinExec/BroadcastNestedLoopJoinExec, andshouldApplySparkToColumnarthenbridges a tagged scan — only in the unsupported-format arms, so the per-format
spark.comet.convert.{csv,json,parquet}.enabledopt-outs still take precedence.The bridge is applied only when:
spark.comet.exec.broadcastExchange.enabledare both on(otherwise the broadcast can't go native, so the bridge would be pointless), and
hasOnlyNativeScans) — i.e. the join canactually become a fully native
CometBroadcastHashJoinExec. Bridging a build side under a jointhat 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 broadcastjoin is Comet, only build a
CometBroadcastExchangeExecfor the reused DPP subquery if thatsubquery's own build is native (
isNativeBuildSide); otherwise fall back to a SparkBroadcastExchangeExec. This prevents an AQE+DPP crash(
Comet execution only takes Arrow Arrays, but got ... OffHeapColumnVector) when the bridgeddimension's separate DPP build copy is row-based. Mirrors the existing non-AQE guard in
CometExecRule.rewriteInSubqueryPlan.How are these changes tested?
CometExecSuite(#6008), covering: a Text build side going native (AQE onand 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;BuildLeftandBroadcastNestedLoopJoinvariants; a non-nativeprobe 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/checkSparkAnswerAndOperatorplus the expected plan shape.CometExecSuitepasses onspark-4.1; the project builds onspark-4.0andspark-4.1.dev/local-ci.sh spark sql_core-1is green, including theDynamicPartitionPruningV1/V2 suites that guard broadcast/DPP reuse (an earlier eager version ofthis change regressed those; the probe-native gate + DPP guard fix it).