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]