This is an automated email from the ASF dual-hosted git repository.

andygrove pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/datafusion-comet.git


The following commit(s) were added to refs/heads/main by this push:
     new 78660f9fcd feat: name incompatible aggregate functions in 
mixed-execution fallback reason (#4750)
78660f9fcd is described below

commit 78660f9fcd9a219ce083178053f0fcb23b0e4076
Author: Andy Grove <[email protected]>
AuthorDate: Wed Jul 8 09:32:52 2026 -0600

    feat: name incompatible aggregate functions in mixed-execution fallback 
reason (#4750)
---
 .../org/apache/comet/serde/QueryPlanSerde.scala    | 30 +++++++++++++++-------
 .../org/apache/spark/sql/comet/operators.scala     | 19 +++++++++-----
 .../approved-plans-v1_4-spark3_5/q70/extended.txt  |  2 +-
 .../approved-plans-v1_4/q10/extended.txt           |  2 +-
 .../approved-plans-v1_4/q35/extended.txt           |  2 +-
 .../approved-plans-v1_4/q45/extended.txt           |  2 +-
 .../approved-plans-v2_7-spark3_5/q70a/extended.txt | 10 ++++----
 .../approved-plans-v2_7/q35/extended.txt           |  2 +-
 8 files changed, 43 insertions(+), 26 deletions(-)

diff --git a/spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala 
b/spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala
index 939d1480ff..91a93d633b 100644
--- a/spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala
+++ b/spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala
@@ -413,15 +413,27 @@ object QueryPlanSerde extends Logging with CometExprShim 
with CometTypeShim {
    * in the other.
    */
   def allAggsSupportMixedExecution(aggExprs: Seq[AggregateExpression]): 
Boolean = {
-    aggExprs.forall { aggExpr =>
-      val fn = aggExpr.aggregateFunction
-      aggrSerdeMap.get(fn.getClass) match {
-        case Some(handler) =>
-          handler
-            .asInstanceOf[CometAggregateExpressionSerde[AggregateFunction]]
-            .supportsMixedPartialFinal
-        case None => false
-      }
+    aggExprs.forall(aggExpr => 
supportsMixedExecution(aggExpr.aggregateFunction))
+  }
+
+  /**
+   * Returns the aggregate functions in the list whose intermediate buffer 
formats are not known
+   * to be compatible between Spark and Comet. These are the functions that 
prevent a Spark Final
+   * aggregate (without a Comet Partial) from running, since the buffer 
produced by one engine
+   * cannot be safely consumed by the other.
+   */
+  def aggsNotSupportingMixedExecution(
+      aggExprs: Seq[AggregateExpression]): Seq[AggregateFunction] = {
+    aggExprs.map(_.aggregateFunction).filterNot(supportsMixedExecution)
+  }
+
+  private def supportsMixedExecution(fn: AggregateFunction): Boolean = {
+    aggrSerdeMap.get(fn.getClass) match {
+      case Some(handler) =>
+        handler
+          .asInstanceOf[CometAggregateExpressionSerde[AggregateFunction]]
+          .supportsMixedPartialFinal
+      case None => false
     }
   }
 
diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala 
b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala
index e4d6b53770..2db1873aa0 100644
--- a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala
+++ b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala
@@ -1538,13 +1538,18 @@ trait CometBaseAggregate {
       return None
     }
 
-    if (sparkFinalMode &&
-      
!QueryPlanSerde.allAggsSupportMixedExecution(aggregate.aggregateExpressions)) {
-      withFallbackReason(
-        aggregate,
-        "Spark Final aggregate without Comet Partial requires compatible " +
-          "intermediate buffer formats")
-      return None
+    if (sparkFinalMode) {
+      val incompatibleAggs =
+        
QueryPlanSerde.aggsNotSupportingMixedExecution(aggregate.aggregateExpressions)
+      if (incompatibleAggs.nonEmpty) {
+        val names = 
incompatibleAggs.map(_.prettyName).distinct.sorted.mkString(", ")
+        withFallbackReason(
+          aggregate,
+          "Spark Final aggregate without Comet Partial requires compatible " +
+            "intermediate buffer formats, but the following aggregate 
function(s) " +
+            s"have incompatible buffers: $names")
+        return None
+      }
     }
 
     // Check if this aggregate has been tagged as unsafe for mixed execution
diff --git 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q70/extended.txt
 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q70/extended.txt
index 14ad77dad4..08df1c04a0 100644
--- 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q70/extended.txt
+++ 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q70/extended.txt
@@ -4,7 +4,7 @@ CometNativeColumnarToRow
       +- CometWindowExec
          +- CometSort
             +- CometColumnarExchange
-               +-  HashAggregate [COMET: Spark Final aggregate without Comet 
Partial requires compatible intermediate buffer formats]
+               +-  HashAggregate [COMET: Spark Final aggregate without Comet 
Partial requires compatible intermediate buffer formats, but the following 
aggregate function(s) have incompatible buffers: sum]
                   +- Exchange
                      +- HashAggregate
                         +- Expand
diff --git 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q10/extended.txt
 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q10/extended.txt
index 07af300183..fcb5e0003d 100644
--- 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q10/extended.txt
+++ 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q10/extended.txt
@@ -1,5 +1,5 @@
 TakeOrderedAndProject
-+-  HashAggregate [COMET: Spark Final aggregate without Comet Partial requires 
compatible intermediate buffer formats]
++-  HashAggregate [COMET: Spark Final aggregate without Comet Partial requires 
compatible intermediate buffer formats, but the following aggregate function(s) 
have incompatible buffers: count]
    +- Exchange
       +- HashAggregate
          +- Project
diff --git 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q35/extended.txt
 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q35/extended.txt
index 07af300183..50e9ebb965 100644
--- 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q35/extended.txt
+++ 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q35/extended.txt
@@ -1,5 +1,5 @@
 TakeOrderedAndProject
-+-  HashAggregate [COMET: Spark Final aggregate without Comet Partial requires 
compatible intermediate buffer formats]
++-  HashAggregate [COMET: Spark Final aggregate without Comet Partial requires 
compatible intermediate buffer formats, but the following aggregate function(s) 
have incompatible buffers: avg, count]
    +- Exchange
       +- HashAggregate
          +- Project
diff --git 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q45/extended.txt
 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q45/extended.txt
index f95c69368f..2d463ef986 100644
--- 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q45/extended.txt
+++ 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4/q45/extended.txt
@@ -1,5 +1,5 @@
 TakeOrderedAndProject
-+-  HashAggregate [COMET: Spark Final aggregate without Comet Partial requires 
compatible intermediate buffer formats]
++-  HashAggregate [COMET: Spark Final aggregate without Comet Partial requires 
compatible intermediate buffer formats, but the following aggregate function(s) 
have incompatible buffers: sum]
    +- Exchange
       +- HashAggregate
          +- Project
diff --git 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q70a/extended.txt
 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q70a/extended.txt
index 10ea854de0..e527d95f46 100644
--- 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q70a/extended.txt
+++ 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q70a/extended.txt
@@ -8,7 +8,7 @@ CometNativeColumnarToRow
                   +- CometColumnarExchange
                      +- HashAggregate
                         +- Union
-                           :-  HashAggregate [COMET: Spark Final aggregate 
without Comet Partial requires compatible intermediate buffer formats]
+                           :-  HashAggregate [COMET: Spark Final aggregate 
without Comet Partial requires compatible intermediate buffer formats, but the 
following aggregate function(s) have incompatible buffers: sum]
                            :  +- Exchange
                            :     +- HashAggregate
                            :        +- Project
@@ -58,10 +58,10 @@ CometNativeColumnarToRow
                            :                                                   
           +- CometProject
                            :                                                   
              +- CometFilter
                            :                                                   
                 +- CometNativeScan parquet spark_catalog.default.date_dim
-                           :-  HashAggregate [COMET: Spark Final aggregate 
without Comet Partial requires compatible intermediate buffer formats]
+                           :-  HashAggregate [COMET: Spark Final aggregate 
without Comet Partial requires compatible intermediate buffer formats, but the 
following aggregate function(s) have incompatible buffers: sum]
                            :  +- Exchange
                            :     +- HashAggregate
-                           :        +-  HashAggregate [COMET: Spark Final 
aggregate without Comet Partial requires compatible intermediate buffer formats]
+                           :        +-  HashAggregate [COMET: Spark Final 
aggregate without Comet Partial requires compatible intermediate buffer 
formats, but the following aggregate function(s) have incompatible buffers: sum]
                            :           +- Exchange
                            :              +- HashAggregate
                            :                 +- Project
@@ -111,10 +111,10 @@ CometNativeColumnarToRow
                            :                                                   
                    +- CometProject
                            :                                                   
                       +- CometFilter
                            :                                                   
                          +- CometNativeScan parquet 
spark_catalog.default.date_dim
-                           +-  HashAggregate [COMET: Spark Final aggregate 
without Comet Partial requires compatible intermediate buffer formats]
+                           +-  HashAggregate [COMET: Spark Final aggregate 
without Comet Partial requires compatible intermediate buffer formats, but the 
following aggregate function(s) have incompatible buffers: sum]
                               +- Exchange
                                  +- HashAggregate
-                                    +-  HashAggregate [COMET: Spark Final 
aggregate without Comet Partial requires compatible intermediate buffer formats]
+                                    +-  HashAggregate [COMET: Spark Final 
aggregate without Comet Partial requires compatible intermediate buffer 
formats, but the following aggregate function(s) have incompatible buffers: sum]
                                        +- Exchange
                                           +- HashAggregate
                                              +- Project
diff --git 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7/q35/extended.txt
 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7/q35/extended.txt
index 07af300183..c340d4a42c 100644
--- 
a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7/q35/extended.txt
+++ 
b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7/q35/extended.txt
@@ -1,5 +1,5 @@
 TakeOrderedAndProject
-+-  HashAggregate [COMET: Spark Final aggregate without Comet Partial requires 
compatible intermediate buffer formats]
++-  HashAggregate [COMET: Spark Final aggregate without Comet Partial requires 
compatible intermediate buffer formats, but the following aggregate function(s) 
have incompatible buffers: avg, count, sum]
    +- Exchange
       +- HashAggregate
          +- Project


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to