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]

Reply via email to