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

zhztheplayer pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git


The following commit(s) were added to refs/heads/main by this push:
     new af05d58e5f [VL] Allow native aggregate offload for a MapType result 
(#13125)
af05d58e5f is described below

commit af05d58e5f63f058a1daf8555a7068b66e008d8c
Author: kevinwilfong <[email protected]>
AuthorDate: Sun Oct 4 06:07:22 2026 -0700

    [VL] Allow native aggregate offload for a MapType result (#13125)
    
    HashAggregateExecBaseTransformer.checkType accepts ArrayType and StructType 
but
    not MapType, so an aggregate carrying a map through its result or buffer
    attributes fails doValidateInternal and runs on the JVM. HashAggregateExec,
    SortAggregateExec and ObjectHashAggregateExec are all offloaded through this
    one transformer, so no route avoids it, and the only trace is a
    GlutenFallbackReporter line reading "Found unsupported data type in 
aggregation
    expression: ...MapType...".
    
    Velox holds a map accumulator for the aggregates that can carry one -- 
arbitrary
    and the spark first / last family use NonNumericArbitrary -- so accept 
MapType
    there.
    
    Widened in the velox transformer rather than the shared base. checkType is a
    statement about what a backend can hold, and the base is inherited by the
    clickhouse and bolt transformers as well; clickhouse already overrides it to
    add StructType, so the per-backend shape is the established one. Only velox 
is
    known to hold a map accumulator here, and a backend that cannot would 
offload
    and fail at execution rather than fall back correctly, so widening the base
    would be asserting something on their behalf that has not been shown.
    
    Spark has no aggregate that builds a map out of non-map input, the way 
presto's
    map_agg does, which is presumably why the gap went unnoticed: the type 
reaches
    checkType only when the data already has a map column and it is carried 
through
    a type-preserving aggregate such as first, last, any_value or max_by, or
    through a UDAF that produces one. The test covers first and last over a map
    column and fails without the change, with the query planned as a vanilla
    SortAggregate.
---
 .../execution/HashAggregateExecTransformer.scala     |  8 ++++++++
 .../execution/VeloxAggregateFunctionsSuite.scala     | 20 ++++++++++++++++++++
 2 files changed, 28 insertions(+)

diff --git 
a/backends-velox/src/main/scala/org/apache/gluten/execution/HashAggregateExecTransformer.scala
 
b/backends-velox/src/main/scala/org/apache/gluten/execution/HashAggregateExecTransformer.scala
index a0c52c7909..3a162aa5c9 100644
--- 
a/backends-velox/src/main/scala/org/apache/gluten/execution/HashAggregateExecTransformer.scala
+++ 
b/backends-velox/src/main/scala/org/apache/gluten/execution/HashAggregateExecTransformer.scala
@@ -65,6 +65,14 @@ abstract class HashAggregateExecTransformer(
     super.output
   }
 
+  // Velox holds a map accumulator for the aggregates that can carry one -- 
arbitrary, and the
+  // spark first / last family, use NonNumericArbitrary. The base transformer 
is shared with the
+  // other backends, which have not been shown to, so widen it here rather 
than there.
+  override protected def checkType(dataType: DataType): Boolean = dataType 
match {
+    case _: MapType => true
+    case other => super.checkType(other)
+  }
+
   override protected def doTransform(context: SubstraitContext): 
TransformContext = {
     val childCtx = child.asInstanceOf[TransformSupport].transform(context)
 
diff --git 
a/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxAggregateFunctionsSuite.scala
 
b/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxAggregateFunctionsSuite.scala
index 6157f356bf..d9c8486154 100644
--- 
a/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxAggregateFunctionsSuite.scala
+++ 
b/backends-velox/src/test/scala/org/apache/gluten/execution/VeloxAggregateFunctionsSuite.scala
@@ -507,6 +507,26 @@ abstract class VeloxAggregateFunctionsSuite extends 
VeloxWholeStageTransformerSu
     }
   }
 
+  test("aggregate over a map column") {
+    // checkType gates the aggregate's result and buffer attributes. Spark has 
no aggregate that
+    // builds a map out of non-map input, so the type only reaches checkType 
when a map column is
+    // carried through a type-preserving aggregate -- first / last / any_value 
/ max_by. Velox
+    // holds a map accumulator for these.
+    withTempView("map_agg_tbl") {
+      Seq((1, Map("a" -> 1.0d)), (1, Map("a" -> 1.0d)), (2, Map("b" -> 2.0d)))
+        .toDF("k", "m")
+        .createOrReplaceTempView("map_agg_tbl")
+
+      // One distinct value per group, so which row first keeps does not 
change the answer.
+      runQueryAndCompare("select k, first(m) from map_agg_tbl group by k") {
+        checkGlutenPlan[HashAggregateExecTransformer]
+      }
+      runQueryAndCompare("select k, last(m) from map_agg_tbl group by k") {
+        checkGlutenPlan[HashAggregateExecTransformer]
+      }
+    }
+  }
+
   test("first") {
     runQueryAndCompare(s"""
                           |select first(l_linenumber), first(l_linenumber, 
true) from lineitem;


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

Reply via email to