From ab4266feb61848c1b084049f877c648c7f11ef1c Mon Sep 17 00:00:00 2001 From: Esteban Zimanyi Date: Mon, 5 Oct 2026 14:51:31 +0200 Subject: [PATCH] Build a typed temporal aggregate without a Spark Column The builder that the spark-sql engine emits for a temporal aggregate of the typed surface made its expression through the public Column API: it applied functions.udaf to new Column(argument) and read the expression back with Column.expr(). Spark 4 removes both, the Column of Spark 4 carrying a ColumnNode rather than an Expression, so the generated TemporalAggregates does not compile against Spark 4. The builder now makes the Catalyst aggregate itself, as the scalar builder of MeosSqlRuntime makes a ScalaUDF: a ScalaAggregator over the argument expressions, the Typed aggregator and the expression encoders of its input and its buffer, turned into an aggregate expression. The constructor and encoderFor have the same signatures in Spark 3.5 and Spark 4, so the generated class compiles against both, and MobilitySpark moves to Spark 4 without a window in which its build and this generator disagree. Witness: with the generator of main, MobilitySpark with Spark 4.0.1 fails to compile TemporalAggregates, "no suitable constructor found for Column(Expression)". Measured: MobilitySpark main regenerated with this generator against MobilityDB 7128c9344a passes its suite, 34 tests and no build warning, with Spark 3.5.1 and with Spark 4.0.1, the temporal aggregates of GeneratedSqlSurfaceTest among them. The codegen tests pass, 106 tests. --- tools/codegen_jvm.py | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/tools/codegen_jvm.py b/tools/codegen_jvm.py index 8841cb6ae..e6fb45d29 100644 --- a/tools/codegen_jvm.py +++ b/tools/codegen_jvm.py @@ -2178,12 +2178,13 @@ def _spark_planning(arity): _SPARK_AGGREGATES = '''package {pkg}; import functions.GeneratedFunctions; -import org.apache.spark.sql.Column; import org.apache.spark.sql.Encoder; import org.apache.spark.sql.Encoders; import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.catalyst.FunctionIdentifier; +import org.apache.spark.sql.catalyst.encoders.package$; import org.apache.spark.sql.catalyst.expressions.Expression; +import org.apache.spark.sql.execution.aggregate.ScalaAggregator; import org.apache.spark.sql.expressions.Aggregator; import org.apache.spark.sql.types.DataType; import org.mobilitydb.spark.generated.TemporalAggregate; @@ -2247,8 +2248,11 @@ def _spark_planning(arity): for (Overload o : overloads) {{ if (o.arg.equals(args.head().dataType())) {{ Aggregator typed = new Typed(o, Encoders.bean(o.out)); - return org.apache.spark.sql.functions.udaf(typed, Encoders.bean(o.in)) - .apply(new Column(args.head())).expr(); + return new ScalaAggregator(args, typed, + package$.MODULE$.encoderFor(Encoders.bean(o.in)), + package$.MODULE$.encoderFor(typed.bufferEncoder()), + true, true, 0, 0, Option.empty()) + .toAggregateExpression(); }} }} }}