From eb8395adc7f011aa5367981577c50ccb21826f32 Mon Sep 17 00:00:00 2001 From: Chang chen Date: Sat, 13 Dec 2025 22:39:32 +0800 Subject: [PATCH 1/5] [GLUTEN-11088][VL] Add compatibility layer for `StructsToJson` across Spark versions and remove Spark-version-specific test limitations --- .../JsonFunctionsValidateSuite.scala | 9 +- .../expression/ExpressionConverter.scala | 216 +++++++++++------- .../objects/InvokeExtractors.scala | 32 +++ .../objects/InvokeExtractors.scala | 32 +++ .../objects/InvokeExtractors.scala | 32 +++ .../objects/InvokeExtractors.scala | 31 +++ .../objects/InvokeExtractors.scala | 46 ++++ 7 files changed, 308 insertions(+), 90 deletions(-) create mode 100644 shims/spark32/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala create mode 100644 shims/spark33/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala create mode 100644 shims/spark34/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala create mode 100644 shims/spark35/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala create mode 100644 shims/spark40/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala diff --git a/backends-velox/src/test/scala/org/apache/gluten/functions/JsonFunctionsValidateSuite.scala b/backends-velox/src/test/scala/org/apache/gluten/functions/JsonFunctionsValidateSuite.scala index 9c80e8cba9d..75f5f2fd431 100644 --- a/backends-velox/src/test/scala/org/apache/gluten/functions/JsonFunctionsValidateSuite.scala +++ b/backends-velox/src/test/scala/org/apache/gluten/functions/JsonFunctionsValidateSuite.scala @@ -59,8 +59,7 @@ class JsonFunctionsValidateSuite extends FunctionsValidateSuite { } } - // TODO: fix on spark-4.0 - testWithMaxSparkVersion("json_array_length", "3.5") { + test("json_array_length") { runQueryAndCompare( s"select *, json_array_length(string_field1) " + s"from datatab limit 5")(checkGlutenPlan[ProjectExecTransformer]) @@ -349,8 +348,7 @@ class JsonFunctionsValidateSuite extends FunctionsValidateSuite { } } - // TODO: fix on spark-4.0 - testWithMaxSparkVersion("json_object_keys", "3.5") { + test("json_object_keys") { withTempPath { path => Seq[String]( @@ -380,8 +378,7 @@ class JsonFunctionsValidateSuite extends FunctionsValidateSuite { } } - // TODO: fix on spark-4.0 - testWithMaxSparkVersion("to_json function", "3.5") { + test("to_json function") { withTable("t") { spark.sql( """ diff --git a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala index c46a9b1077b..18aa65d4d3d 100644 --- a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala +++ b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala @@ -26,7 +26,7 @@ import org.apache.spark.{SPARK_REVISION, SPARK_VERSION_SHORT} import org.apache.spark.internal.Logging import org.apache.spark.sql.catalyst.SQLConfHelper import org.apache.spark.sql.catalyst.expressions.{StringTrimBoth, _} -import org.apache.spark.sql.catalyst.expressions.objects.StaticInvoke +import org.apache.spark.sql.catalyst.expressions.objects.{Invoke, StaticInvoke, StructsToJsonInvoke} import org.apache.spark.sql.catalyst.optimizer.NormalizeNaNAndZero import org.apache.spark.sql.execution.ScalarSubquery import org.apache.spark.sql.hive.HiveUDFTransformer @@ -144,89 +144,110 @@ object ExpressionConverter extends SQLConfHelper with Logging { DecimalArithmeticExpressionTransformer(substraitName, leftChild, rightChild, resultType, b) } + // Mapping for Iceberg static invoke functions + private val icebergStaticInvokeMap = Map( + "BucketFunction" -> ExpressionNames.BUCKET, + "TruncateFunction" -> ExpressionNames.TRUNCATE, + "YearsFunction" -> ExpressionNames.YEARS, + "MonthsFunction" -> ExpressionNames.MONTHS, + "DaysFunction" -> ExpressionNames.DAYS, + "HoursFunction" -> ExpressionNames.HOURS + ) + + // Mapping for other static invoke functions + private val staticInvokeMap = Map( + "varcharTypeWriteSideCheck" -> ExpressionNames.VARCHAR_TYPE_WRITE_SIDE_CHECK, + "charTypeWriteSideCheck" -> ExpressionNames.CHAR_TYPE_WRITE_SIDE_CHECK, + "readSidePadding" -> ExpressionNames.READ_SIDE_PADDING, + "lengthOfJsonArray" -> ExpressionNames.JSON_ARRAY_LENGTH, + "jsonObjectKeys" -> ExpressionNames.JSON_OBJECT_KEYS + ) + private def replaceStaticInvokeWithExpressionTransformer( i: StaticInvoke, attributeSeq: Seq[Attribute], expressionsMap: Map[Class[_], String]): ExpressionTransformer = { + + val objName = i.staticObject.getName + val funcName = i.functionName + + def doTransform(child: Expression): ExpressionTransformer = + replaceWithExpressionTransformer0(child, attributeSeq, expressionsMap) + def validateAndTransform( exprName: String, childTransformers: => Seq[ExpressionTransformer]): ExpressionTransformer = { if (!BackendsApiManager.getValidatorApiInstance.doExprValidate(exprName, i)) { throw new GlutenNotSupportException( - s"Not supported to map current ${i.getClass} call on function: ${i.functionName}.") + s"Not supported to map current ${i.getClass} call on function: $funcName.") } GenericExpressionTransformer(exprName, childTransformers, i) } - i.functionName match { - case "encode" | "decode" if i.objectName.endsWith("UrlCodec") => - validateAndTransform( - "url_" + i.functionName, - Seq(replaceWithExpressionTransformer0(i.arguments.head, attributeSeq, expressionsMap)) - ) + // Try to match Iceberg static invoke first + val icebergTransformer: Option[ExpressionTransformer] = + if (funcName == "invoke") { + icebergStaticInvokeMap.collectFirst { + case (func, name) if objName.startsWith("org.apache.iceberg.spark.functions." + func) => + GenericExpressionTransformer(name, i.arguments.map(doTransform), i) + } + } else { + None + } - case "isLuhnNumber" => - validateAndTransform( - ExpressionNames.LUHN_CHECK, - Seq(replaceWithExpressionTransformer0(i.arguments.head, attributeSeq, expressionsMap)) - ) + icebergTransformer.getOrElse { + funcName match { + case "isLuhnNumber" => + validateAndTransform(ExpressionNames.LUHN_CHECK, Seq(doTransform(i.arguments.head))) - case "encode" | "decode" if i.objectName.endsWith("Base64") => - if (!BackendsApiManager.getValidatorApiInstance.doExprValidate(ExpressionNames.BASE64, i)) { - throw new GlutenNotSupportException( - s"Not supported to map current ${i.getClass} call on function: ${i.functionName}.") - } - BackendsApiManager.getSparkPlanExecApiInstance.genBase64StaticInvokeTransformer( - ExpressionNames.BASE64, - replaceWithExpressionTransformer0(i.arguments.head, attributeSeq, expressionsMap), - i - ) + case fn @ ("encode" | "decode") if objName.endsWith("UrlCodec") => + validateAndTransform("url_" + fn, Seq(doTransform(i.arguments.head))) - case fn - if i.objectName.endsWith("CharVarcharCodegenUtils") && Set( - "varcharTypeWriteSideCheck", - "charTypeWriteSideCheck", - "readSidePadding").contains(fn) => - val exprName = fn match { - case "varcharTypeWriteSideCheck" => ExpressionNames.VARCHAR_TYPE_WRITE_SIDE_CHECK - case "charTypeWriteSideCheck" => ExpressionNames.CHAR_TYPE_WRITE_SIDE_CHECK - case "readSidePadding" => ExpressionNames.READ_SIDE_PADDING - } - validateAndTransform( - exprName, - i.arguments.map(replaceWithExpressionTransformer0(_, attributeSeq, expressionsMap)) - ) + case "encode" | "decode" if objName.endsWith("Base64") => + if ( + !BackendsApiManager.getValidatorApiInstance.doExprValidate(ExpressionNames.BASE64, i) + ) { + throw new GlutenNotSupportException( + s"Not supported to map current ${i.getClass} call on function: $funcName.") + } + BackendsApiManager.getSparkPlanExecApiInstance.genBase64StaticInvokeTransformer( + ExpressionNames.BASE64, + doTransform(i.arguments.head), + i + ) - case _ => - throw new GlutenNotSupportException( - s"Not supported to transform StaticInvoke with object: ${i.staticObject.getName}, " + - s"function: ${i.functionName}") + case fn if staticInvokeMap.contains(fn) => + validateAndTransform(staticInvokeMap(fn), i.arguments.map(doTransform)) + + case _ => + throw new GlutenNotSupportException( + s"Not supported to transform StaticInvoke with object: $objName, function: $funcName") + } } } - private def replaceIcebergStaticInvoke( - s: StaticInvoke, + private def replaceInvokeWithExpressionTransformer( + invoke: Invoke, attributeSeq: Seq[Attribute], expressionsMap: Map[Class[_], String]): ExpressionTransformer = { - val invokeMap = Map( - "BucketFunction" -> ExpressionNames.BUCKET, - "TruncateFunction" -> ExpressionNames.TRUNCATE, - "YearsFunction" -> ExpressionNames.YEARS, - "MonthsFunction" -> ExpressionNames.MONTHS, - "DaysFunction" -> ExpressionNames.DAYS, - "HoursFunction" -> ExpressionNames.HOURS - ) - val objName = s.staticObject.getName - val transformer = invokeMap.find { - case (func, _) => objName.startsWith("org.apache.iceberg.spark.functions." + func) - } - if (transformer.isEmpty) { - throw new GlutenNotSupportException(s"Not supported staticInvoke call object: $objName") + + // Pattern matching for different Invoke types + invoke match { + // StructsToJson evaluator + case StructsToJsonInvoke(options, child, timeZoneId) => + val toJsonExpr = StructsToJson(options, child, timeZoneId) + val substraitExprName = getAndCheckSubstraitName(toJsonExpr, expressionsMap) + BackendsApiManager.getSparkPlanExecApiInstance.genToJsonTransformer( + substraitExprName, + replaceWithExpressionTransformer0(child, attributeSeq, expressionsMap), + toJsonExpr + ) + + // Unsupported invoke + case _ => + throw new GlutenNotSupportException( + s"Not supported to transform Invoke with function: $invoke") } - GenericExpressionTransformer( - transformer.get._2, - s.arguments.map(replaceWithExpressionTransformer0(_, attributeSeq, expressionsMap)), - s) } private def replaceWithExpressionTransformer0( @@ -237,32 +258,59 @@ object ExpressionConverter extends SQLConfHelper with Logging { s"replaceWithExpressionTransformer expr: $expr class: ${expr.getClass} " + s"name: ${expr.prettyName}") - expr match { - case p: PythonUDF => - return replacePythonUDFWithExpressionTransformer(p, attributeSeq, expressionsMap) - case s: ScalaUDF => - return replaceScalaUDFWithExpressionTransformer(s, attributeSeq, expressionsMap) - case _ if HiveUDFTransformer.isHiveUDF(expr) => - return BackendsApiManager.getSparkPlanExecApiInstance.genHiveUDFTransformer( - expr, - attributeSeq) - case i: StaticInvoke - if i.functionName == "invoke" && i.staticObject.getName.startsWith( - "org.apache.iceberg.spark.functions.") => - return replaceIcebergStaticInvoke(i, attributeSeq, expressionsMap) - case i: StaticInvoke => - return replaceStaticInvokeWithExpressionTransformer(i, attributeSeq, expressionsMap) - case _ => + tryTransformWithoutExpressionMapping(expr, attributeSeq, expressionsMap).getOrElse { + val substraitExprName: String = getAndCheckSubstraitName(expr, expressionsMap) + val backendConverted = BackendsApiManager.getSparkPlanExecApiInstance + .extraExpressionConverter(substraitExprName, expr, attributeSeq) + + backendConverted.getOrElse( + transformExpression(expr, attributeSeq, expressionsMap, substraitExprName)) } + } - val substraitExprName: String = getAndCheckSubstraitName(expr, expressionsMap) - val backendConverted = BackendsApiManager.getSparkPlanExecApiInstance.extraExpressionConverter( - substraitExprName, - expr, - attributeSeq) - if (backendConverted.isDefined) { - return backendConverted.get + /** + * Transform expressions that don't have direct expression class mapping in expressionsMap. This + * handles special cases like UDFs (Python, Scala, Hive), StaticInvoke, and Invoke expressions, + * where the transformation logic is based on runtime information rather than expression class + * type. + * + * @param expr + * The expression to transform + * @param attributeSeq + * The sequence of attributes for binding + * @param expressionsMap + * The expression class to substrait name mapping (not used for these cases) + * @return + * Some(ExpressionTransformer) if the expression matches one of these special cases, None + * otherwise + */ + private def tryTransformWithoutExpressionMapping( + expr: Expression, + attributeSeq: Seq[Attribute], + expressionsMap: Map[Class[_], String]): Option[ExpressionTransformer] = { + Option { + expr match { + case pythonUDF: PythonUDF => + replacePythonUDFWithExpressionTransformer(pythonUDF, attributeSeq, expressionsMap) + case scalaUDF: ScalaUDF => + replaceScalaUDFWithExpressionTransformer(scalaUDF, attributeSeq, expressionsMap) + case _ if HiveUDFTransformer.isHiveUDF(expr) => + BackendsApiManager.getSparkPlanExecApiInstance.genHiveUDFTransformer(expr, attributeSeq) + case staticInvoke: StaticInvoke => + replaceStaticInvokeWithExpressionTransformer(staticInvoke, attributeSeq, expressionsMap) + case invoke: Invoke => + replaceInvokeWithExpressionTransformer(invoke, attributeSeq, expressionsMap) + case _ => + null + } } + } + + private def transformExpression( + expr: Expression, + attributeSeq: Seq[Attribute], + expressionsMap: Map[Class[_], String], + substraitExprName: String): ExpressionTransformer = { expr match { case c: CreateArray => val children = diff --git a/shims/spark32/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala b/shims/spark32/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala new file mode 100644 index 00000000000..24687a2e62d --- /dev/null +++ b/shims/spark32/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala @@ -0,0 +1,32 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.spark.sql.catalyst.expressions.objects + +import org.apache.spark.sql.catalyst.expressions.Expression + +/** + * Extractors for Invoke expressions to ensure compatibility across different Spark versions. + * + * For Spark 3.2, StructsToJson is not replaced with Invoke expressions, + * so this extractor returns None to maintain API compatibility with other versions. + */ +object StructsToJsonInvoke { + def unapply(expr: Expression): Option[(Map[String, String], Expression, Option[String])] = { + None + } +} + diff --git a/shims/spark33/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala b/shims/spark33/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala new file mode 100644 index 00000000000..5f0ed024af8 --- /dev/null +++ b/shims/spark33/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala @@ -0,0 +1,32 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.spark.sql.catalyst.expressions.objects + +import org.apache.spark.sql.catalyst.expressions.Expression + +/** + * Extractors for Invoke expressions to ensure compatibility across different Spark versions. + * + * For Spark 3.3, StructsToJson is not replaced with Invoke expressions, + * so this extractor returns None to maintain API compatibility with other versions. + */ +object StructsToJsonInvoke { + def unapply(expr: Expression): Option[(Map[String, String], Expression, Option[String])] = { + None + } +} + diff --git a/shims/spark34/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala b/shims/spark34/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala new file mode 100644 index 00000000000..5b39ae2fb08 --- /dev/null +++ b/shims/spark34/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala @@ -0,0 +1,32 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.spark.sql.catalyst.expressions.objects + +import org.apache.spark.sql.catalyst.expressions.Expression + +/** + * Extractors for Invoke expressions to ensure compatibility across different Spark versions. + * + * For Spark 3.4, StructsToJson is not replaced with Invoke expressions, + * so this extractor returns None to maintain API compatibility with other versions. + */ +object StructsToJsonInvoke { + def unapply(expr: Expression): Option[(Map[String, String], Expression, Option[String])] = { + None + } +} + diff --git a/shims/spark35/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala b/shims/spark35/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala new file mode 100644 index 00000000000..6bfbd77eec3 --- /dev/null +++ b/shims/spark35/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala @@ -0,0 +1,31 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.spark.sql.catalyst.expressions.objects + +import org.apache.spark.sql.catalyst.expressions.Expression + +/** + * Extractors for Invoke expressions to ensure compatibility across different Spark versions. + * + * For Spark 3.5, StructsToJson is not replaced with Invoke expressions, so this extractor returns + * None to maintain API compatibility with other versions. + */ +object StructsToJsonInvoke { + def unapply(expr: Expression): Option[(Map[String, String], Expression, Option[String])] = { + None + } +} diff --git a/shims/spark40/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala b/shims/spark40/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala new file mode 100644 index 00000000000..8abebd0d2a8 --- /dev/null +++ b/shims/spark40/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala @@ -0,0 +1,46 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.spark.sql.catalyst.expressions.objects + +import org.apache.spark.sql.catalyst.expressions.{Expression, Literal} +import org.apache.spark.sql.catalyst.expressions.json.StructsToJsonEvaluator + +/** + * Extractors for Invoke expressions to ensure compatibility across different Spark versions. + * + * Since Spark 4.0, StructsToJson has been replaced with Invoke expressions using + * StructsToJsonEvaluator. This extractor provides a unified interface to extract evaluator options, + * child expression, and timeZoneId from the Invoke pattern. + */ +object StructsToJsonInvoke { + def unapply(expr: Expression): Option[(Map[String, String], Expression, Option[String])] = { + expr match { + case Invoke( + Literal(evaluator: StructsToJsonEvaluator, _), + "evaluate", + _, + Seq(child), + _, + _, + _, + _) => + Some((evaluator.options, child, evaluator.timeZoneId)) + case _ => + None + } + } +} From 7578e3cf61ee21367a9ca7031869b8ad9f07572a Mon Sep 17 00:00:00 2001 From: Chang chen Date: Mon, 15 Dec 2025 10:31:45 +0800 Subject: [PATCH 2/5] Chore: Fix style issue --- .../sql/catalyst/expressions/objects/InvokeExtractors.scala | 5 ++--- .../sql/catalyst/expressions/objects/InvokeExtractors.scala | 5 ++--- 2 files changed, 4 insertions(+), 6 deletions(-) diff --git a/shims/spark32/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala b/shims/spark32/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala index 24687a2e62d..a8b215b8892 100644 --- a/shims/spark32/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala +++ b/shims/spark32/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala @@ -21,12 +21,11 @@ import org.apache.spark.sql.catalyst.expressions.Expression /** * Extractors for Invoke expressions to ensure compatibility across different Spark versions. * - * For Spark 3.2, StructsToJson is not replaced with Invoke expressions, - * so this extractor returns None to maintain API compatibility with other versions. + * For Spark 3.2, StructsToJson is not replaced with Invoke expressions, so this extractor returns + * None to maintain API compatibility with other versions. */ object StructsToJsonInvoke { def unapply(expr: Expression): Option[(Map[String, String], Expression, Option[String])] = { None } } - diff --git a/shims/spark34/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala b/shims/spark34/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala index 5b39ae2fb08..c4982946bfa 100644 --- a/shims/spark34/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala +++ b/shims/spark34/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala @@ -21,12 +21,11 @@ import org.apache.spark.sql.catalyst.expressions.Expression /** * Extractors for Invoke expressions to ensure compatibility across different Spark versions. * - * For Spark 3.4, StructsToJson is not replaced with Invoke expressions, - * so this extractor returns None to maintain API compatibility with other versions. + * For Spark 3.4, StructsToJson is not replaced with Invoke expressions, so this extractor returns + * None to maintain API compatibility with other versions. */ object StructsToJsonInvoke { def unapply(expr: Expression): Option[(Map[String, String], Expression, Option[String])] = { None } } - From 2da8c36bc8803a850a20770fd436de32e4dd39ee Mon Sep 17 00:00:00 2001 From: Chang chen Date: Mon, 15 Dec 2025 10:54:49 +0800 Subject: [PATCH 3/5] Chore: Fix style issue --- .../sql/catalyst/expressions/objects/InvokeExtractors.scala | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/shims/spark33/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala b/shims/spark33/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala index 5f0ed024af8..95372c082cb 100644 --- a/shims/spark33/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala +++ b/shims/spark33/src/main/scala/org/apache/spark/sql/catalyst/expressions/objects/InvokeExtractors.scala @@ -21,12 +21,11 @@ import org.apache.spark.sql.catalyst.expressions.Expression /** * Extractors for Invoke expressions to ensure compatibility across different Spark versions. * - * For Spark 3.3, StructsToJson is not replaced with Invoke expressions, - * so this extractor returns None to maintain API compatibility with other versions. + * For Spark 3.3, StructsToJson is not replaced with Invoke expressions, so this extractor returns + * None to maintain API compatibility with other versions. */ object StructsToJsonInvoke { def unapply(expr: Expression): Option[(Map[String, String], Expression, Option[String])] = { None } } - From bd3d05211c58cfa433e4f1bd401d0a1421a20374 Mon Sep 17 00:00:00 2001 From: Chang chen Date: Mon, 15 Dec 2025 11:36:51 +0800 Subject: [PATCH 4/5] Fix: using `staticObject.objectName` --- .../org/apache/gluten/expression/ExpressionConverter.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala index 18aa65d4d3d..b3b8756f8cd 100644 --- a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala +++ b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala @@ -168,7 +168,7 @@ object ExpressionConverter extends SQLConfHelper with Logging { attributeSeq: Seq[Attribute], expressionsMap: Map[Class[_], String]): ExpressionTransformer = { - val objName = i.staticObject.getName + val objName = i.objectName val funcName = i.functionName def doTransform(child: Expression): ExpressionTransformer = From 00a7b7847d81849fed841dbae3d8b5552c6303da Mon Sep 17 00:00:00 2001 From: Chang chen Date: Mon, 15 Dec 2025 12:02:48 +0800 Subject: [PATCH 5/5] Refactor: refine Base64 encode/decode handling with streamlined validation logic --- .../apache/gluten/expression/ExpressionConverter.scala | 10 +++------- 1 file changed, 3 insertions(+), 7 deletions(-) diff --git a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala index b3b8756f8cd..52f6d31d1d1 100644 --- a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala +++ b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionConverter.scala @@ -203,13 +203,9 @@ object ExpressionConverter extends SQLConfHelper with Logging { case fn @ ("encode" | "decode") if objName.endsWith("UrlCodec") => validateAndTransform("url_" + fn, Seq(doTransform(i.arguments.head))) - case "encode" | "decode" if objName.endsWith("Base64") => - if ( - !BackendsApiManager.getValidatorApiInstance.doExprValidate(ExpressionNames.BASE64, i) - ) { - throw new GlutenNotSupportException( - s"Not supported to map current ${i.getClass} call on function: $funcName.") - } + case fn @ ("encode" | "decode") + if objName.endsWith("Base64") && BackendsApiManager.getValidatorApiInstance + .doExprValidate(ExpressionNames.BASE64, i) => BackendsApiManager.getSparkPlanExecApiInstance.genBase64StaticInvokeTransformer( ExpressionNames.BASE64, doTransform(i.arguments.head),