kazuyukitanimura commented on code in PR #4565:
URL: https://github.com/apache/datafusion-comet/pull/4565#discussion_r4168965742


##########
spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:
##########


Review Comment:
   Following up on @sunchao's point about SortExec. I think the gap is wider 
than the distinct case. Spark always plans a Sort between a sort-aggregate 
Final and its  ▎ exchange, so findPartialAggInPlan never reaches the Partial 
for any sort aggregate. That means the #1389 tagging pass is effectively off 
for SortAggregateExec.
    Main has since added revertUnsafePartialAggregates and 
hasUnrepairedNativeBuffer as a backstop, but both match only 
CometHashAggregateExec and stop at a
    CometSortExec. Could the rebase widen all three to CometBaseAggregateExec 
and walk through SortExec / CometSortExec? It would also be great to run the 
existing  ▎ "mixed engine collect_list" tests a second time with 
useObjectHashAggregateExec=false, so both halves of the split are covered for 
sort aggregate



##########
spark/src/main/scala/org/apache/spark/sql/comet/operators.scala:
##########
@@ -2051,29 +2043,72 @@ object CometObjectHashAggregateExec
   }
 }
 
-case class CometHashAggregateExec(
-    override val nativeOp: Operator,
-    override val originalPlan: SparkPlan,
-    override val output: Seq[Attribute],
-    groupingExpressions: Seq[NamedExpression],
-    aggregateExpressions: Seq[AggregateExpression],
-    resultExpressions: Seq[NamedExpression],
-    input: Seq[Attribute],
-    child: SparkPlan,
-    override val serializedPlanOpt: SerializedPlan)
+object CometSortAggregateExec
+    extends CometOperatorSerde[SortAggregateExec]
+    with CometBaseAggregate {
+
+  override def enabledConfig: Option[ConfigEntry[Boolean]] = Some(
+    CometConf.COMET_EXEC_AGGREGATE_ENABLED)
+
+  override def getSupportLevel(op: SortAggregateExec): SupportLevel =
+    baseAggregateSupportLevel(op)
+
+  override def convert(
+      aggregate: SortAggregateExec,
+      builder: Operator.Builder,
+      childOp: OperatorOuterClass.Operator*): 
Option[OperatorOuterClass.Operator] = {
+
+    // SortAggregate is planned for TypedImperativeAggregate functions whose 
intermediate
+    // buffer formats differ between Spark and Comet (same risk as 
ObjectHashAggregate).
+    // Require Comet shuffle so a Partial->Final pair never spans the 
JVM/native boundary.
+    if (!isCometShuffleEnabled(aggregate.conf)) {
+      return None
+    }
+
+    doConvert(aggregate, builder, childOp: _*)
+  }
+
+  override def createExec(nativeOp: Operator, op: SortAggregateExec): 
CometNativeExec = {
+    // The native AggregateExec auto-detects Sorted input mode from the 
child's output ordering

Review Comment:
    I checked @sunchao's array-key repro against DataFusion 55.1 and the cause 
is in GroupValuesColumn::vectorized_intern. A row whose hash matches an 
existing bucket but has a different value is placed later, in 
scalarized_intern_remaining, after all the new-hash rows in that batch. A NULL 
list and an empty list hash the same, so [] ends up after [1]. So first-seen 
emission isn't something we can rely on for nested or multi-column keys. Since 
outputOrdering advertises grouping-key order and Spark has already dropped the 
downstream sorts, could CometSortAggregateExec only convert when the ordering 
is produced inside the same native block (for example, the child is a 
CometSortExec)? Inserting a native sort or falling back otherwise would also 
work. A test with a cached, sorted relation grouped by an array key would lock 
this in.



##########
dev/diffs/4.1.3.diff:
##########
@@ -2010,26 +2010,31 @@ index 46ed8fdfd21..4585dbab5b8 100644
  
    private def checkWindowGroupLimits(query: String, count: Int): Unit = {
 diff --git 
a/sql/core/src/test/scala/org/apache/spark/sql/execution/ReplaceHashWithSortAggSuite.scala
 
b/sql/core/src/test/scala/org/apache/spark/sql/execution/ReplaceHashWithSortAggSuite.scala
-index 47679ed7865..9ffbaecb98e 100644
+index 47679ed7865..cbcabe00992 100644

Review Comment:
   we now need 4.2.0.diff 



-- 
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