Skip to content

AQE reuses one exchange for scans with different dynamic pruning filters and drops rows #6264

Description

@dwsmith1983

Describe the bug

With AQE on, two exchanges that differ only in the dynamic partition pruning filter of a scan below them can compare equal, and AQE replaces the second with a ReusedExchange of the first. The consumer of the second exchange then reads the other branch's rows, and those rows are silently lost from the result.

CometScanUtils.filterUnusedDynamicPruningExpressions drops DynamicPruningExpression(TrueLiteral) like Spark does. It also drops any pruning filter whose subquery is still the adaptive placeholder (SubqueryAdaptiveBroadcastExec or CometSubqueryAdaptiveBroadcastExec). Spark's FileSourceScanExec keeps that filter, and its canonical form includes the dimension side's plan, so IN (1999 dates) and IN (2000 dates) stay distinct.

AQE canonicalizes a query stage from its exchange as it was before the stage optimizer rules ran (ExchangeQueryStageExec._canonicalized). That is before CometPlanAdaptiveDynamicPruningFilters converts the placeholder. So any exchange above that stage sees the scan with no pruning filter at all. Two such parent exchanges, one over each scan, get the same canonical form, and the stage cache hands out the first one twice.

A coalesced AQEShuffleRead between the parent and the child stage hides the problem, because its partition specs carry per-partition data sizes that differ between the two branches. AQE adds no such read when coalescing is off, or when coalescing would not merge any partitions, which is what happens with large shuffles. This is probably why TPC-DS q64 returns 0 rows only sometimes, and only at large scale factors, in #6133: the two cross_sales references build the same store_sales join for 1999 and 2000. That link has not been confirmed at SF1000.

This is separate from the q5 failure in #6133, which comes from CometNativeScanExec.outputPartitioning running the placeholder subquery.

Steps to reproduce

Reproduced on main (605051a) with -Pspark-3.5 and -Pspark-4.1, with the whole plan running in Comet (native Parquet scan and Comet shuffle enabled).

spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "false")

spark.range(1000)
  .selectExpr("cast(id % 10 as int) as store_id", "cast(id % 13 as int) as item", "cast(id as int) as amount")
  .write.partitionBy("store_id").parquet("/tmp/fact")
spark.range(10)
  .selectExpr("cast(id as int) as store_id",
    "case when id = 1 then 'a' when id in (2, 3, 4) then 'b' else 'c' end as grp")
  .write.parquet("/tmp/dim")
spark.read.parquet("/tmp/fact").createOrReplaceTempView("fact")
spark.read.parquet("/tmp/dim").createOrReplaceTempView("dim")

spark.sql("""
  WITH x AS (SELECT store_id, item, sum(amount) AS s FROM fact GROUP BY store_id, item),
  y AS (SELECT store_id, s % 7 AS b, count(*) AS c FROM x GROUP BY store_id, s % 7)
  SELECT y.store_id, y.b, y.c FROM y JOIN dim d ON y.store_id = d.store_id WHERE d.grp = 'a'
  UNION ALL
  SELECT y.store_id, y.b, y.c FROM y JOIN dim d ON y.store_id = d.store_id WHERE d.grp = 'b'
""").collect().length

Spark returns 28 rows. Comet returns 7 on Spark 3.5 and 21 on Spark 4.1, depending on which branch's stage is created first. The final plan has a ReusedExchange under the second branch pointing at the first branch's exchange, and only one of the two pruned scans runs.

The same thing happens when the parent is a broadcast, for example a broadcast of a sort merge join over the scan's shuffle stage, which is the q64 shape. With coalescing on, these small tables give correct results because the coalesced reads keep the two parents apart.

Expected behavior

Same rows as Spark. Two scans with different dynamic pruning filters should never share an exchange.

Additional context

The extra stripping was added in #4112 so that scans whose pruning filter is still unconverted could reuse each other. Matching Spark's rule and removing the stale filter copy inside CometNativeScanExec.originalPlan from its canonical form fixes the reproduction while keeping the SPARK-32509 reuse test passing. CometScanExec and CometIcebergNativeScanExec call the same helper.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions