andygrove commented on code in PR #4565:
URL: https://github.com/apache/datafusion-comet/pull/4565#discussion_r4169488087
##########
spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:
##########
Review Comment:
Agreed, it affected every sort aggregate, not just the distinct case. The
merge widened `restoreSparkPartial`, `revertUnsafePartialAggregates` and
`hasUnrepairedNativeBuffer` to `CometBaseAggregateExec` (41c7baf26d), and
c506925525 makes all three passes walk through the sort: `SortExec` in
`findPartialAggInPlan`, and `SortExec` or `CometSortExec` in the two
post-conversion passes. The mixed-engine collect_list tests now also run with
`useObjectHashAggregateExec=false` and count every Comet aggregate wrapper
(c66e94e998). Without the traversal fix, the Comet partial + Spark final
variant fails with an NPE in Spark's Final.
##########
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:
Yes, that's the cause. I went with the native sort: sort aggregates now set
`ordered_by_grouping_keys` on the HashAggregate proto, and the native planner
sorts the aggregate output on the grouping columns whenever the AggregateExec's
output ordering doesn't already satisfy them (1c16c7d15e). When the sort is in
the same native block, DataFusion picks the Sorted mode and nothing is added.
The cached, sorted relation with array keys is now a `CometInMemoryCacheSuite`
test, and it returned `[1]` before `[]` without the fix.
--
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]