comphead commented on issue #4412:
URL:
https://github.com/apache/datafusion-comet/issues/4412#issuecomment-5895745889
Repro for both halves of this issue, including the wrong results reported in
the comment above. Two different `BaseAggregateExec` patterns in Spark are
involved, and only the global-aggregate one changes results.
**Setup:** `local[2]`, upstream Spark 4.1.3, AQE on, default Comet configs
plus `CometShuffleManager` and off-heap. Comet built on main `5ca149928f`
(2026-09-21), plus unrelated round-robin shuffle commits. main still has
`CometHashAggregateExec extends CometUnaryExec`, so I expect the same on
current main, but I haven't rerun there.
```scala
val data = spark.range(0, 20000, 1, 4)
.selectExpr("cast(id as string) as content_id", "cast(id % 13 as int) as
genre_id")
data.write.parquet("/tmp/i4412/t1")
data.write.parquet("/tmp/i4412/t2")
spark.range(0, 1000, 1, 2).selectExpr("id as x", "if(id % 2 = 0, 'a', 'b')
as p")
.write.partitionBy("p").parquet("/tmp/i4412/tp")
Seq("t1", "t2", "tp").foreach(t =>
spark.read.parquet(s"/tmp/i4412/$t").createOrReplaceTempView(t))
```
### 1. Grouped aggregate: lost optimization, correct results
`AQEPropagateEmptyRelation.getEstimatedRowCount` reads the child stage's row
count only through `LogicalQueryStage(_, agg: BaseAggregateExec) if
agg.groupingExpressions.nonEmpty`.
```sql
SELECT genre_id, count(*) FROM t1 WHERE content_id = 'no_match' GROUP BY
genre_id;
SELECT content_id, genre_id FROM t1 EXCEPT SELECT content_id, genre_id FROM
t2;
```
| Query | Spark | Comet |
|---|---|---|
| `GROUP BY`, filter matches nothing | 0 rows, final plan `EmptyRelation`, 1
job | 0 rows, `CometHashAggregate [Final]` over an empty `ShuffleQueryStage`, 2
jobs |
| `EXCEPT` of identical inputs | 0 rows, `EmptyRelation`, 2 jobs | 0 rows,
same pattern, 3 jobs |
The same shows up in a production event log on Spark 3.4.3. An `EXCEPT` of
two identical inputs collapsed to an empty `LocalTableScan` on Spark, while
Comet ran one more job for the final aggregate over the empty shuffle.
### 2. Global aggregate: wrong results
SPARK-44040 added a second pattern, in `LogicalQueryStage.computeStats`. A
`BaseAggregateExec` with no grouping keys over a stage with 0 rows reports
`rowCount = 1`, because a global aggregate always returns a row. A Comet final
aggregate doesn't match, so its stage reports 0 rows. `OptimizeOneRowPlan` then
sees `maxRows = 0` for a union of such stages and removes the `DISTINCT` above
it.
This is the `SPARK-44040` test from `AdaptiveQueryExecSuite`, which is
`IgnoreComet(#4412)` in every `dev/diffs` file since #4861:
```scala
val emptyDf = spark.range(1).where("false")
val a1 = emptyDf.agg(sum("id").as("id")).withColumn("name", lit("df1"))
val a2 = emptyDf.agg(sum("id").as("id")).withColumn("name", lit("df2"))
a1.union(a2).select("id").distinct().collect()
```
It also happens on a fully native Parquet plan with default configs, when
partition pruning leaves no files:
```sql
SELECT DISTINCT s FROM (
SELECT sum(x) AS s FROM tp WHERE p = 'none'
UNION ALL
SELECT sum(x) AS s FROM tp WHERE p = 'nope');
```
| Query | Spark | Comet |
|---|---|---|
| SPARK-44040 with `sum` | `[null]` | `[null], [null]` |
| SPARK-44040 with `count` | `[0]` | `[0]` |
| SPARK-44040 with `count`, `spark.comet.exec.localTableScan.enabled=true` |
`[0]` | `[0], [0]` |
| Parquet, partitions pruned to no files | `[null]` | `[null], [null]` |
In every wrong case the Comet final plan is `CometProject <- CometUnion <-
CometHashAggregate [Final]`, with the `DISTINCT` gone. `count` comes out right
only in the mixed case, because #4242 keeps a Spark partial `count` from
feeding a native final. Once the partial is native too, `count` fails the same
way.
### Notes on the description
- The two `SPARK-35442` tests it lists are not ignored in `dev/diffs` today,
since #4374 didn't merge. The test tagged with this issue is `SPARK-44040`.
- "Query results are correct under Comet" holds for grouped aggregates only.
- Option C (`CometHashAggregateExec extends BaseAggregateExec`) should cover
both patterns, though I haven't tested it. Option B as described targets the
grouped one. The wrong result comes from `LogicalQueryStage.computeStats`, so B
would have to handle that path too.
--
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]