From d7a30ba1fbb6560edff420d3e597089ce0a88a26 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Thu, 20 Aug 2026 12:35:49 +0530 Subject: [PATCH 01/14] [SPARK-50698][SQL] Refactor CreateUserDefinedFunction command to extend from UnaryRunnableCommand ### What changes were proposed in this pull request? This PR refactors `CreateUserDefinedFunctionCommand` to extend from `UnaryRunnableCommand`, taking `child: LogicalPlan`. It follows up on #49126 (SPARK-48730) by removing the duplicate `CreateUserDefinedFunction` Catalyst logical plan from `v2Commands.scala` and unifying the command structure. ### Why are the changes needed? To simplify and unify the logical command abstractions for SQL UDFs introduced in #49126. ### Does this PR introduce any user-facing change? No. ### How was this patch tested? - Updated unit tests in `CreateSQLFunctionParserSuite`. --- .../analysis/ApplyDefaultCollation.scala | 13 +----- .../catalyst/analysis/ResolveCatalogs.scala | 6 --- .../catalyst/plans/logical/v2Commands.scala | 22 +-------- .../analysis/ResolveSessionCatalog.scala | 22 ++------- .../catalyst/parser/SqlStatementCodes.scala | 3 +- .../spark/sql/execution/SparkSqlParser.scala | 5 ++- .../command/CreateSQLFunctionCommand.scala | 34 +++++++------- .../CreateUserDefinedFunctionCommand.scala | 45 +++++++++++++++++-- .../CreateSQLFunctionParserSuite.scala | 9 ++-- 9 files changed, 72 insertions(+), 87 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala index ff43b3668839b..02dc94b019388 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala @@ -21,7 +21,7 @@ import scala.util.control.NonFatal import org.apache.spark.SparkException import org.apache.spark.sql.catalyst.expressions.{AttributeReference, Cast, DefaultStringProducingExpression, Expression, Literal, SubqueryExpression} -import org.apache.spark.sql.catalyst.plans.logical.{AddColumns, AlterColumns, AlterColumnSpec, AlterViewAs, ColumnDefinition, CreateTable, CreateTableAsSelect, CreateTempView, CreateUserDefinedFunction, CreateView, LogicalPlan, QualifiedColType, ReplaceColumns, ReplaceTable, ReplaceTableAsSelect, TableSpec, V2CreateTablePlan} +import org.apache.spark.sql.catalyst.plans.logical.{AddColumns, AlterColumns, AlterColumnSpec, AlterViewAs, ColumnDefinition, CreateTable, CreateTableAsSelect, CreateTempView, CreateView, LogicalPlan, QualifiedColType, ReplaceColumns, ReplaceTable, ReplaceTableAsSelect, TableSpec, V2CreateTablePlan} import org.apache.spark.sql.catalyst.rules.Rule import org.apache.spark.sql.catalyst.trees.CurrentOrigin import org.apache.spark.sql.catalyst.types.DataTypeUtils.{areSameBaseType, isDefaultStringCharOrVarcharType, replaceDefaultStringCharAndVarcharTypes} @@ -206,17 +206,6 @@ object ApplyDefaultCollation extends Rule[LogicalPlan] { newCreateView.copyTagsFrom(createView) newCreateView - case createUserDefinedFunction@CreateUserDefinedFunction( - ResolvedIdentifier(catalog: SupportsNamespaces, identifier), - _, _, _, _, _, collation, _, _, _, _, _, _) if collation.isEmpty => - val newCreateUserDefinedFunction = - CurrentOrigin.withOrigin(createUserDefinedFunction.origin) { - createUserDefinedFunction.copy( - collation = getCollationFromSchemaMetadata(catalog, identifier.namespace())) - } - newCreateUserDefinedFunction.copyTagsFrom(createUserDefinedFunction) - newCreateUserDefinedFunction - case other => other } diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala index 6fc196774a048..6419427169109 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala @@ -91,12 +91,6 @@ class ResolveCatalogs(val catalogManager: CatalogManager) throw QueryCompilationErrors.operationNotAllowedOnBuiltinFunctionError( "CREATE", nameParts.last) - case CreateUserDefinedFunction(UnresolvedIdentifier(nameParts, _), - _, _, _, _, _, _, _, _, _, _, _, _) - if isSystemBuiltinName(nameParts) => - throw QueryCompilationErrors.operationNotAllowedOnBuiltinFunctionError( - "CREATE", nameParts.last) - case DropFunction(UnresolvedIdentifier(nameParts, _), _) if isSystemBuiltinName(nameParts) => throw QueryCompilationErrors.operationNotAllowedOnBuiltinFunctionError( diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala index b816016a3ec84..f5acafc70abb2 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala @@ -21,7 +21,7 @@ import org.apache.spark.{SparkException, SparkIllegalArgumentException, SparkUns import org.apache.spark.sql.AnalysisException import org.apache.spark.sql.catalyst.analysis.{AnalysisContext, AssignmentUtils, EliminateSubqueryAliases, FieldName, NamedRelation, PartitionSpec, ResolvedIdentifier, ResolvedProcedure, ResolveSchemaEvolution, TypeCheckResult, UnresolvedAttribute, UnresolvedException, UnresolvedProcedure, ViewSchemaMode} import org.apache.spark.sql.catalyst.analysis.TypeCheckResult.{DataTypeMismatch, TypeCheckSuccess} -import org.apache.spark.sql.catalyst.catalog.{FunctionResource, RoutineLanguage} +import org.apache.spark.sql.catalyst.catalog.FunctionResource import org.apache.spark.sql.catalyst.catalog.CatalogTypes.TablePartitionSpec import org.apache.spark.sql.catalyst.expressions._ import org.apache.spark.sql.catalyst.plans.DescribeCommandSchema @@ -1629,26 +1629,6 @@ case class CreateFunction( copy(child = newChild) } -/** - * The logical plan of the CREATE FUNCTION command for SQL Functions. - */ -case class CreateUserDefinedFunction( - child: LogicalPlan, - inputParamText: Option[String], - returnTypeText: String, - exprText: Option[String], - queryText: Option[String], - comment: Option[String], - collation: Option[String], - isDeterministic: Option[Boolean], - containsSQL: Option[Boolean], - language: RoutineLanguage, - isTableFunc: Boolean, - ignoreIfExists: Boolean, - replace: Boolean) extends UnaryCommand { - override protected def withNewChildInternal(newChild: LogicalPlan): CreateUserDefinedFunction = - copy(child = newChild) -} /** * The logical plan of the DROP FUNCTION command. diff --git a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala index e6fc6d8d862ce..c35adc27d23b1 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala @@ -698,25 +698,11 @@ class ResolveSessionCatalog(val catalogManager: CatalogManager) case CreateFunction(ResolvedIdentifier(catalog, _), _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) - case c @ CreateUserDefinedFunction( + case c @ CreateSQLFunctionCommand( CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _) => - CreateUserDefinedFunctionCommand( - FunctionIdentifier(ident.table, ident.database, ident.catalog), - c.inputParamText, - c.returnTypeText, - c.exprText, - c.queryText, - c.comment, - c.collation, - c.isDeterministic, - c.containsSQL, - c.language, - c.isTableFunc, - isTemp = false, - c.ignoreIfExists, - c.replace) - - case CreateUserDefinedFunction( + c + + case CreateSQLFunctionCommand( ResolvedIdentifier(catalog, _), _, _, _, _, _, _, _, _, _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/parser/SqlStatementCodes.scala b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/parser/SqlStatementCodes.scala index 24cfdfcd2eb49..e7b17a8827dbd 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/parser/SqlStatementCodes.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/parser/SqlStatementCodes.scala @@ -146,8 +146,7 @@ object SqlStatementCodes { case _: SetCatalogAndNamespace | _: SetNamespaceCommand => SetSchema case _: SetCatalogCommand => SetCatalog case _: TruncateTable => TruncateTable - case _: CreateFunction | _: CreateFunctionCommand | - _: CreateUserDefinedFunction | _: CreateUserDefinedFunctionCommand => + case _: CreateFunction | _: CreateFunctionCommand | _: CreateUserDefinedFunctionCommand => CreateRoutine case _: DropFunction | _: DropFunctionCommand => DropRoutine case _: UnresolvedExecuteImmediate => ExecuteImmediate diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala index 5070a40259e37..bf2655f755e0f 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala @@ -1041,7 +1041,7 @@ class SparkSqlAstBuilder extends AstBuilder { withIdentClause(ctx.identifierReference(), functionIdentifier => { if (ctx.TEMPORARY == null) { - CreateUserDefinedFunction( + CreateUserDefinedFunctionCommand( UnresolvedIdentifier(functionIdentifier), inputParamText, returnTypeText, @@ -1053,6 +1053,7 @@ class SparkSqlAstBuilder extends AstBuilder { containsSQL, language, isTableFunc, + isTemp = false, ctx.EXISTS != null, ctx.REPLACE != null) } else { @@ -1064,7 +1065,7 @@ class SparkSqlAstBuilder extends AstBuilder { // Extract the actual function name, handling session qualification val funcName = extractTempFunctionName(functionIdentifier, ctx) CreateUserDefinedFunctionCommand( - FunctionIdentifier(funcName), + UnresolvedIdentifier(Seq(funcName)), inputParamText, returnTypeText, exprText, diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala index a597087085b42..809123f87749b 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala @@ -20,7 +20,7 @@ package org.apache.spark.sql.execution.command import org.apache.spark.SparkException import org.apache.spark.sql.{AnalysisException, Row, SparkSession} import org.apache.spark.sql.catalyst.FunctionIdentifier -import org.apache.spark.sql.catalyst.analysis.{withPosition, Analyzer, SQLFunctionExpression, SQLFunctionNode, SQLScalarFunction, SQLTableFunction, UnresolvedAlias, UnresolvedAttribute, UnresolvedFunction, UnresolvedRelation, UnresolvedTableValuedFunction} +import org.apache.spark.sql.catalyst.analysis.{withPosition, Analyzer, ResolvedIdentifier, SQLFunctionExpression, SQLFunctionNode, SQLScalarFunction, SQLTableFunction, UnresolvedAlias, UnresolvedAttribute, UnresolvedFunction, UnresolvedIdentifier, UnresolvedRelation, UnresolvedTableValuedFunction} import org.apache.spark.sql.catalyst.catalog.{SessionCatalog, SQLFunction, UserDefinedFunction, UserDefinedFunctionErrors} import org.apache.spark.sql.catalyst.catalog.UserDefinedFunction._ import org.apache.spark.sql.catalyst.expressions.{Alias, Cast, Expression, Generator, LateralSubquery, Literal, ScalarSubquery, SubqueryExpression, WindowExpression} @@ -34,24 +34,8 @@ import org.apache.spark.sql.errors.QueryCompilationErrors import org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand._ import org.apache.spark.sql.types.{DataType, MetadataBuilder, StructField, StructType} -/** - * The DDL command that creates a SQL function. - * For example: - * {{{ - * CREATE [OR REPLACE] [TEMPORARY] FUNCTION [IF NOT EXISTS] [db_name.]function_name - * ([param_name param_type [COMMENT param_comment], ...]) - * RETURNS {ret_type | TABLE (ret_name ret_type [COMMENT ret_comment], ...])} - * [function_properties] function_body; - * - * function_properties: - * [NOT] DETERMINISTIC | COMMENT function_comment | [ CONTAINS SQL | READS SQL DATA ] - * - * function_body: - * RETURN {expression | TABLE ( query )} - * }}} - */ case class CreateSQLFunctionCommand( - name: FunctionIdentifier, + child: LogicalPlan, inputParamText: Option[String], returnTypeText: String, exprText: Option[String], @@ -68,7 +52,21 @@ case class CreateSQLFunctionCommand( import SQLFunction._ + lazy val name: FunctionIdentifier = child match { + case ResolvedIdentifier(c, ident) => + FunctionIdentifier(ident.name(), ident.namespace().headOption) + case u: UnresolvedIdentifier => + FunctionIdentifier(u.nameParts.last, u.nameParts.dropRight(1).lastOption) + case _ => + throw SparkException.internalError( + s"Unexpected child plan in CreateSQLFunctionCommand: $child") + } + + override protected def withNewChildInternal( + newChild: LogicalPlan): CreateSQLFunctionCommand = copy(child = newChild) + override def run(sparkSession: SparkSession): Seq[Row] = { + val parser = sparkSession.sessionState.sqlParser val analyzer = sparkSession.sessionState.analyzer val catalog = sparkSession.sessionState.catalog diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateUserDefinedFunctionCommand.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateUserDefinedFunctionCommand.scala index f65c7c91251a1..2b336604948f0 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateUserDefinedFunctionCommand.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateUserDefinedFunctionCommand.scala @@ -25,11 +25,14 @@ import org.apache.spark.sql.catalyst.catalog.{LanguageSQL, RoutineLanguage, User import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types.StructType +import org.apache.spark.sql.catalyst.analysis.UnresolvedIdentifier +import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan + /** * The base class for CreateUserDefinedFunctionCommand */ abstract class CreateUserDefinedFunctionCommand - extends LeafRunnableCommand with CapturesConfig + extends UnaryRunnableCommand with CapturesConfig object CreateUserDefinedFunctionCommand { @@ -40,7 +43,7 @@ object CreateUserDefinedFunctionCommand { */ // scalastyle:off argcount def apply( - name: FunctionIdentifier, + child: LogicalPlan, inputParamText: Option[String], returnTypeText: String, exprText: Option[String], @@ -62,7 +65,7 @@ object CreateUserDefinedFunctionCommand { language match { case LanguageSQL => CreateSQLFunctionCommand( - name, + child, inputParamText, returnTypeText, exprText, @@ -80,6 +83,42 @@ object CreateUserDefinedFunctionCommand { throw UserDefinedFunctionErrors.unsupportedUserDefinedFunction(other) } } + // scalastyle:off argcount + def apply( + name: FunctionIdentifier, + inputParamText: Option[String], + returnTypeText: String, + exprText: Option[String], + queryText: Option[String], + comment: Option[String], + collation: Option[String], + isDeterministic: Option[Boolean], + containsSQL: Option[Boolean], + language: RoutineLanguage, + isTableFunc: Boolean, + isTemp: Boolean, + ignoreIfExists: Boolean, + replace: Boolean + ): CreateUserDefinedFunctionCommand = { + // scalastyle:on argcount + val nameParts = name.database.toSeq :+ name.funcName + apply( + UnresolvedIdentifier(nameParts), + inputParamText, + returnTypeText, + exprText, + queryText, + comment, + collation, + isDeterministic, + containsSQL, + language, + isTableFunc, + isTemp, + ignoreIfExists, + replace) + } + /** * Check whether the function parameters contain duplicated column names. diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala index 56316f43f8dfe..b9b5fc1a0a7ff 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala @@ -18,10 +18,8 @@ package org.apache.spark.sql.execution.command import org.apache.spark.sql.AnalysisException -import org.apache.spark.sql.catalyst.FunctionIdentifier import org.apache.spark.sql.catalyst.analysis.{AnalysisTest, UnresolvedIdentifier} import org.apache.spark.sql.catalyst.catalog.LanguageSQL -import org.apache.spark.sql.catalyst.plans.logical.CreateUserDefinedFunction import org.apache.spark.sql.execution.SparkSqlParser class CreateSQLFunctionParserSuite extends AnalysisTest { @@ -49,9 +47,9 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL: Option[Boolean] = None, isTableFunc: Boolean = false, ignoreIfExists: Boolean = false, - replace: Boolean = false): CreateUserDefinedFunction = { + replace: Boolean = false): CreateUserDefinedFunctionCommand = { // scalastyle:on argcount - CreateUserDefinedFunction( + CreateUserDefinedFunctionCommand( UnresolvedIdentifier(nameParts), inputParamText = inputParamText, returnTypeText = returnTypeText, @@ -63,6 +61,7 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL = containsSQL, language = LanguageSQL, isTableFunc = isTableFunc, + isTemp = false, ignoreIfExists = ignoreIfExists, replace = replace) } @@ -82,7 +81,7 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { replace: Boolean = false): CreateSQLFunctionCommand = { // scalastyle:on argcount CreateSQLFunctionCommand( - FunctionIdentifier(name), + UnresolvedIdentifier(Seq(name)), inputParamText = inputParamText, returnTypeText = returnTypeText, exprText = exprText, From 1732e217e49d352d505ca052a38b125971b4de62 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Thu, 20 Aug 2026 12:57:07 +0530 Subject: [PATCH 02/14] [SPARK-50698][SQL] Fix scalastyle import ordering --- .../apache/spark/sql/catalyst/plans/logical/v2Commands.scala | 2 +- .../execution/command/CreateUserDefinedFunctionCommand.scala | 5 ++--- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala index f5acafc70abb2..805b8c6ed55cd 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala @@ -21,8 +21,8 @@ import org.apache.spark.{SparkException, SparkIllegalArgumentException, SparkUns import org.apache.spark.sql.AnalysisException import org.apache.spark.sql.catalyst.analysis.{AnalysisContext, AssignmentUtils, EliminateSubqueryAliases, FieldName, NamedRelation, PartitionSpec, ResolvedIdentifier, ResolvedProcedure, ResolveSchemaEvolution, TypeCheckResult, UnresolvedAttribute, UnresolvedException, UnresolvedProcedure, ViewSchemaMode} import org.apache.spark.sql.catalyst.analysis.TypeCheckResult.{DataTypeMismatch, TypeCheckSuccess} -import org.apache.spark.sql.catalyst.catalog.FunctionResource import org.apache.spark.sql.catalyst.catalog.CatalogTypes.TablePartitionSpec +import org.apache.spark.sql.catalyst.catalog.FunctionResource import org.apache.spark.sql.catalyst.expressions._ import org.apache.spark.sql.catalyst.plans.DescribeCommandSchema import org.apache.spark.sql.catalyst.trees.BinaryLike diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateUserDefinedFunctionCommand.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateUserDefinedFunctionCommand.scala index 2b336604948f0..a73d01e1c9145 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateUserDefinedFunctionCommand.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateUserDefinedFunctionCommand.scala @@ -21,13 +21,12 @@ import java.util.Locale import org.apache.spark.sql.AnalysisException import org.apache.spark.sql.catalyst.{CapturesConfig, FunctionIdentifier} +import org.apache.spark.sql.catalyst.analysis.UnresolvedIdentifier import org.apache.spark.sql.catalyst.catalog.{LanguageSQL, RoutineLanguage, UserDefinedFunctionErrors} +import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types.StructType -import org.apache.spark.sql.catalyst.analysis.UnresolvedIdentifier -import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan - /** * The base class for CreateUserDefinedFunctionCommand */ From 5d198b49fe560113c98b96f631f2a0060355944b Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Thu, 20 Aug 2026 16:34:52 +0530 Subject: [PATCH 03/14] [SPARK-50698][SQL] Refactor CreateUserDefinedFunctionCommand to extend from UnaryRunnableCommand --- .../analysis/ResolveSessionCatalog.scala | 24 +++++++++++++++---- .../CreateSQLFunctionParserSuite.scala | 6 ++--- 2 files changed, 22 insertions(+), 8 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala index c35adc27d23b1..e8358172b1c9a 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala @@ -698,11 +698,25 @@ class ResolveSessionCatalog(val catalogManager: CatalogManager) case CreateFunction(ResolvedIdentifier(catalog, _), _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) - case c @ CreateSQLFunctionCommand( - CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _) => - c - - case CreateSQLFunctionCommand( + case c @ CreateUserDefinedFunction( + child @ CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _) => + CreateUserDefinedFunctionCommand( + child, + c.inputParamText, + c.returnTypeText, + c.exprText, + c.queryText, + c.comment, + c.collation, + c.isDeterministic, + c.containsSQL, + c.language, + c.isTableFunc, + isTemp = false, + c.ignoreIfExists, + c.replace) + + case CreateUserDefinedFunction( ResolvedIdentifier(catalog, _), _, _, _, _, _, _, _, _, _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala index b9b5fc1a0a7ff..dc83d40abdd8d 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala @@ -20,6 +20,7 @@ package org.apache.spark.sql.execution.command import org.apache.spark.sql.AnalysisException import org.apache.spark.sql.catalyst.analysis.{AnalysisTest, UnresolvedIdentifier} import org.apache.spark.sql.catalyst.catalog.LanguageSQL +import org.apache.spark.sql.catalyst.plans.logical.CreateUserDefinedFunction import org.apache.spark.sql.execution.SparkSqlParser class CreateSQLFunctionParserSuite extends AnalysisTest { @@ -47,9 +48,9 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL: Option[Boolean] = None, isTableFunc: Boolean = false, ignoreIfExists: Boolean = false, - replace: Boolean = false): CreateUserDefinedFunctionCommand = { + replace: Boolean = false): CreateUserDefinedFunction = { // scalastyle:on argcount - CreateUserDefinedFunctionCommand( + CreateUserDefinedFunction( UnresolvedIdentifier(nameParts), inputParamText = inputParamText, returnTypeText = returnTypeText, @@ -61,7 +62,6 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL = containsSQL, language = LanguageSQL, isTableFunc = isTableFunc, - isTemp = false, ignoreIfExists = ignoreIfExists, replace = replace) } From 011628211864590bf3b8236433d5b81b0cfb0eb6 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Thu, 20 Aug 2026 18:10:59 +0530 Subject: [PATCH 04/14] [SPARK-50698][SQL] Add MiMa excludes for CreateUserDefinedFunctionCommand refactoring --- project/MimaExcludes.scala | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/project/MimaExcludes.scala b/project/MimaExcludes.scala index c27ef26959edf..840649b38239c 100644 --- a/project/MimaExcludes.scala +++ b/project/MimaExcludes.scala @@ -70,7 +70,12 @@ object MimaExcludes { // [SPARK-57987] Add desc field to the SQL REST API Node case class ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.status.api.v1.sql.Node.apply"), ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.status.api.v1.sql.Node.copy"), - ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.status.api.v1.sql.Node$") + ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.status.api.v1.sql.Node$"), + // [SPARK-50698][SQL] Refactor CreateUserDefinedFunctionCommand to extend from UnaryRunnableCommand + ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand"), + ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand.apply"), + ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.sql.execution.command.CreateSQLFunctionCommand"), + ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.apply") ) // Exclude rules for 4.2.x from 4.1.0 From 784b5062c2187546ff45c2851d9d13a1d107ea93 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Fri, 21 Aug 2026 00:03:08 +0530 Subject: [PATCH 05/14] [SPARK-50698][SQL] Trigger fresh GitHub Actions CI run From 54d84363fde7c64cfc85f895c271e848429d4f7f Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Fri, 21 Aug 2026 01:17:29 +0530 Subject: [PATCH 06/14] [SPARK-50698][SQL] Add MiMa excludes for removed CreateUserDefinedFunction logical plan node --- project/MimaExcludes.scala | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/project/MimaExcludes.scala b/project/MimaExcludes.scala index 840649b38239c..6ad9255dc1bab 100644 --- a/project/MimaExcludes.scala +++ b/project/MimaExcludes.scala @@ -75,7 +75,11 @@ object MimaExcludes { ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand"), ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand.apply"), ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.sql.execution.command.CreateSQLFunctionCommand"), - ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.apply") + ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.apply"), + ProblemFilters.exclude[MissingClassProblem]( + "org.apache.spark.sql.catalyst.plans.logical.CreateUserDefinedFunction"), + ProblemFilters.exclude[MissingClassProblem]( + "org.apache.spark.sql.catalyst.plans.logical.CreateUserDefinedFunction$") ) // Exclude rules for 4.2.x from 4.1.0 From 3b0d542d88cf7c5032a74f99aa07e01e74b93aa1 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Fri, 21 Aug 2026 01:41:58 +0530 Subject: [PATCH 07/14] [SPARK-50698][SQL] Fix pattern matching and suite types for CreateSQLFunctionCommand --- .../analysis/ResolveSessionCatalog.scala | 24 ++++--------------- .../CreateSQLFunctionParserSuite.scala | 9 ++++--- 2 files changed, 9 insertions(+), 24 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala index e8358172b1c9a..c35adc27d23b1 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala @@ -698,25 +698,11 @@ class ResolveSessionCatalog(val catalogManager: CatalogManager) case CreateFunction(ResolvedIdentifier(catalog, _), _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) - case c @ CreateUserDefinedFunction( - child @ CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _) => - CreateUserDefinedFunctionCommand( - child, - c.inputParamText, - c.returnTypeText, - c.exprText, - c.queryText, - c.comment, - c.collation, - c.isDeterministic, - c.containsSQL, - c.language, - c.isTableFunc, - isTemp = false, - c.ignoreIfExists, - c.replace) - - case CreateUserDefinedFunction( + case c @ CreateSQLFunctionCommand( + CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _) => + c + + case CreateSQLFunctionCommand( ResolvedIdentifier(catalog, _), _, _, _, _, _, _, _, _, _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala index dc83d40abdd8d..24c5f297f0a04 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala @@ -19,8 +19,7 @@ package org.apache.spark.sql.execution.command import org.apache.spark.sql.AnalysisException import org.apache.spark.sql.catalyst.analysis.{AnalysisTest, UnresolvedIdentifier} -import org.apache.spark.sql.catalyst.catalog.LanguageSQL -import org.apache.spark.sql.catalyst.plans.logical.CreateUserDefinedFunction +import org.apache.spark.sql.execution.command.CreateSQLFunctionCommand import org.apache.spark.sql.execution.SparkSqlParser class CreateSQLFunctionParserSuite extends AnalysisTest { @@ -48,9 +47,9 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL: Option[Boolean] = None, isTableFunc: Boolean = false, ignoreIfExists: Boolean = false, - replace: Boolean = false): CreateUserDefinedFunction = { + replace: Boolean = false): CreateSQLFunctionCommand = { // scalastyle:on argcount - CreateUserDefinedFunction( + CreateSQLFunctionCommand( UnresolvedIdentifier(nameParts), inputParamText = inputParamText, returnTypeText = returnTypeText, @@ -60,8 +59,8 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { collation = None, isDeterministic = isDeterministic, containsSQL = containsSQL, - language = LanguageSQL, isTableFunc = isTableFunc, + isTemp = false, ignoreIfExists = ignoreIfExists, replace = replace) } From 1057313be63758691a97acad53aa06b04c13a144 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Fri, 21 Aug 2026 16:50:29 +0530 Subject: [PATCH 08/14] [SPARK-50698][SQL] Fix resolution of CreateUserDefinedFunction logical plan node --- .../analysis/ApplyDefaultCollation.scala | 13 ++++++++++- .../catalyst/analysis/ResolveCatalogs.scala | 8 +++++++ .../catalyst/plans/logical/v2Commands.scala | 21 ++++++++++++++++++ .../analysis/ResolveSessionCatalog.scala | 22 +++++++++++++++---- .../spark/sql/execution/SparkSqlParser.scala | 9 +++----- .../CreateSQLFunctionParserSuite.scala | 14 ++++++------ 6 files changed, 69 insertions(+), 18 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala index 02dc94b019388..103dd45ed3062 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala @@ -21,7 +21,7 @@ import scala.util.control.NonFatal import org.apache.spark.SparkException import org.apache.spark.sql.catalyst.expressions.{AttributeReference, Cast, DefaultStringProducingExpression, Expression, Literal, SubqueryExpression} -import org.apache.spark.sql.catalyst.plans.logical.{AddColumns, AlterColumns, AlterColumnSpec, AlterViewAs, ColumnDefinition, CreateTable, CreateTableAsSelect, CreateTempView, CreateView, LogicalPlan, QualifiedColType, ReplaceColumns, ReplaceTable, ReplaceTableAsSelect, TableSpec, V2CreateTablePlan} +import org.apache.spark.sql.catalyst.plans.logical.{AddColumns, AlterColumns, AlterColumnSpec, AlterViewAs, ColumnDefinition, CreateTable, CreateTableAsSelect, CreateTempView, CreateUserDefinedFunction, CreateView, LogicalPlan, QualifiedColType, ReplaceColumns, ReplaceTable, ReplaceTableAsSelect, TableSpec, V2CreateTablePlan} import org.apache.spark.sql.catalyst.rules.Rule import org.apache.spark.sql.catalyst.trees.CurrentOrigin import org.apache.spark.sql.catalyst.types.DataTypeUtils.{areSameBaseType, isDefaultStringCharOrVarcharType, replaceDefaultStringCharAndVarcharTypes} @@ -206,6 +206,17 @@ object ApplyDefaultCollation extends Rule[LogicalPlan] { newCreateView.copyTagsFrom(createView) newCreateView + case createUserDefinedFunction@CreateUserDefinedFunction(ResolvedIdentifier( + catalog: SupportsNamespaces, identifier), _, _, _, _, _, _, _, _, _, _, _, _) + if createUserDefinedFunction.collation.isEmpty => + val newCreateUserDefinedFunction = + CurrentOrigin.withOrigin(createUserDefinedFunction.origin) { + createUserDefinedFunction.copy( + collation = getCollationFromSchemaMetadata(catalog, identifier.namespace())) + } + newCreateUserDefinedFunction.copyTagsFrom(createUserDefinedFunction) + newCreateUserDefinedFunction + case other => other } diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala index 6419427169109..45001dad6b08c 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala @@ -91,6 +91,14 @@ class ResolveCatalogs(val catalogManager: CatalogManager) throw QueryCompilationErrors.operationNotAllowedOnBuiltinFunctionError( "CREATE", nameParts.last) + case c @ CreateUserDefinedFunction( + u @ UnresolvedIdentifier(nameParts, _), _, _, _, _, _, _, _, _, _, _, _, _) => + if (isSystemBuiltinName(nameParts)) { + throw QueryCompilationErrors.operationNotAllowedOnBuiltinFunctionError( + "CREATE", nameParts.last) + } + c.copy(child = resolveFunctionIdentifier(nameParts, u.origin)) + case DropFunction(UnresolvedIdentifier(nameParts, _), _) if isSystemBuiltinName(nameParts) => throw QueryCompilationErrors.operationNotAllowedOnBuiltinFunctionError( diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala index 805b8c6ed55cd..32ab97c94ecf1 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala @@ -1629,6 +1629,27 @@ case class CreateFunction( copy(child = newChild) } +/** + * The logical plan of the CREATE FUNCTION command for SQL Functions. + */ +case class CreateUserDefinedFunction( + child: LogicalPlan, + inputParamText: Option[String], + returnTypeText: String, + exprText: Option[String], + queryText: Option[String], + comment: Option[String], + collation: Option[String], + isDeterministic: Option[Boolean], + containsSQL: Option[Boolean], + language: org.apache.spark.sql.catalyst.catalog.RoutineLanguage, + isTableFunc: Boolean, + ignoreIfExists: Boolean, + replace: Boolean) extends UnaryCommand { + override protected def withNewChildInternal(newChild: LogicalPlan): CreateUserDefinedFunction = + copy(child = newChild) +} + /** * The logical plan of the DROP FUNCTION command. diff --git a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala index c35adc27d23b1..fe52e4613b7e3 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala @@ -698,11 +698,25 @@ class ResolveSessionCatalog(val catalogManager: CatalogManager) case CreateFunction(ResolvedIdentifier(catalog, _), _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) - case c @ CreateSQLFunctionCommand( + case c @ CreateUserDefinedFunction( CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _) => - c - - case CreateSQLFunctionCommand( + CreateUserDefinedFunctionCommand( + child = c.child, + inputParamText = c.inputParamText, + returnTypeText = c.returnTypeText, + exprText = c.exprText, + queryText = c.queryText, + comment = c.comment, + collation = c.collation, + isDeterministic = c.isDeterministic, + containsSQL = c.containsSQL, + language = c.language, + isTableFunc = c.isTableFunc, + isTemp = false, + ignoreIfExists = c.ignoreIfExists, + replace = c.replace) + + case CreateUserDefinedFunction( ResolvedIdentifier(catalog, _), _, _, _, _, _, _, _, _, _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala index bf2655f755e0f..3a7e0bc5b3087 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala @@ -1041,7 +1041,7 @@ class SparkSqlAstBuilder extends AstBuilder { withIdentClause(ctx.identifierReference(), functionIdentifier => { if (ctx.TEMPORARY == null) { - CreateUserDefinedFunctionCommand( + CreateUserDefinedFunction( UnresolvedIdentifier(functionIdentifier), inputParamText, returnTypeText, @@ -1053,7 +1053,6 @@ class SparkSqlAstBuilder extends AstBuilder { containsSQL, language, isTableFunc, - isTemp = false, ctx.EXISTS != null, ctx.REPLACE != null) } else { @@ -1064,7 +1063,7 @@ class SparkSqlAstBuilder extends AstBuilder { // Extract the actual function name, handling session qualification val funcName = extractTempFunctionName(functionIdentifier, ctx) - CreateUserDefinedFunctionCommand( + CreateUserDefinedFunction( UnresolvedIdentifier(Seq(funcName)), inputParamText, returnTypeText, @@ -1076,10 +1075,8 @@ class SparkSqlAstBuilder extends AstBuilder { containsSQL, language, isTableFunc, - isTemp = true, ctx.EXISTS != null, - ctx.REPLACE != null - ) + ctx.REPLACE != null) } }) } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala index 24c5f297f0a04..c6e7992e7d60d 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala @@ -19,7 +19,7 @@ package org.apache.spark.sql.execution.command import org.apache.spark.sql.AnalysisException import org.apache.spark.sql.catalyst.analysis.{AnalysisTest, UnresolvedIdentifier} -import org.apache.spark.sql.execution.command.CreateSQLFunctionCommand +import org.apache.spark.sql.catalyst.plans.logical.CreateUserDefinedFunction import org.apache.spark.sql.execution.SparkSqlParser class CreateSQLFunctionParserSuite extends AnalysisTest { @@ -47,9 +47,9 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL: Option[Boolean] = None, isTableFunc: Boolean = false, ignoreIfExists: Boolean = false, - replace: Boolean = false): CreateSQLFunctionCommand = { + replace: Boolean = false): CreateUserDefinedFunction = { // scalastyle:on argcount - CreateSQLFunctionCommand( + CreateUserDefinedFunction( UnresolvedIdentifier(nameParts), inputParamText = inputParamText, returnTypeText = returnTypeText, @@ -59,8 +59,8 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { collation = None, isDeterministic = isDeterministic, containsSQL = containsSQL, + language = org.apache.spark.sql.catalyst.catalog.LanguageSQL, isTableFunc = isTableFunc, - isTemp = false, ignoreIfExists = ignoreIfExists, replace = replace) } @@ -77,9 +77,9 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL: Option[Boolean] = None, isTableFunc: Boolean = false, ignoreIfExists: Boolean = false, - replace: Boolean = false): CreateSQLFunctionCommand = { + replace: Boolean = false): CreateUserDefinedFunction = { // scalastyle:on argcount - CreateSQLFunctionCommand( + CreateUserDefinedFunction( UnresolvedIdentifier(Seq(name)), inputParamText = inputParamText, returnTypeText = returnTypeText, @@ -89,8 +89,8 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { collation = None, isDeterministic = isDeterministic, containsSQL = containsSQL, + language = org.apache.spark.sql.catalyst.catalog.LanguageSQL, isTableFunc = isTableFunc, - isTemp = true, ignoreIfExists = ignoreIfExists, replace = replace) } From f29c708a041ff581cec8a4b2d1a1ba00a0f73ce0 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Thu, 27 Aug 2026 19:50:30 +0530 Subject: [PATCH 09/14] [SPARK-50698][SQL] Add Scaladoc for CreateSQLFunctionCommand --- .../command/CreateSQLFunctionCommand.scala | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala index 809123f87749b..c424c2f6f9bf0 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala @@ -34,6 +34,22 @@ import org.apache.spark.sql.errors.QueryCompilationErrors import org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand._ import org.apache.spark.sql.types.{DataType, MetadataBuilder, StructField, StructType} +/** + * The DDL command that creates a SQL function. + * For example: + * {{{ + * CREATE [OR REPLACE] [TEMPORARY] FUNCTION [IF NOT EXISTS] [db_name.]function_name + * ([param_name param_type [COMMENT param_comment], ...]) + * RETURNS {ret_type | TABLE (ret_name ret_type [COMMENT ret_comment], ...])} + * [function_properties] function_body; + * + * function_properties: + * [NOT] DETERMINISTIC | COMMENT function_comment | [ CONTAINS SQL | READS SQL DATA ] + * + * function_body: + * RETURN {expression | TABLE ( query )} + * }}} + */ case class CreateSQLFunctionCommand( child: LogicalPlan, inputParamText: Option[String], From aa51097e03427a3d1d255305027e47950c70a9d9 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Thu, 27 Aug 2026 22:53:44 +0530 Subject: [PATCH 10/14] [SPARK-50698][SQL] Preserve isTemp flag in CreateUserDefinedFunction logical plan node --- .../spark/sql/catalyst/analysis/ApplyDefaultCollation.scala | 2 +- .../spark/sql/catalyst/analysis/ResolveCatalogs.scala | 2 +- .../spark/sql/catalyst/plans/logical/v2Commands.scala | 1 + .../spark/sql/catalyst/analysis/ResolveSessionCatalog.scala | 6 +++--- .../org/apache/spark/sql/execution/SparkSqlParser.scala | 2 ++ .../execution/command/CreateSQLFunctionParserSuite.scala | 2 ++ 6 files changed, 10 insertions(+), 5 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala index 103dd45ed3062..8ebf0f09eb357 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ApplyDefaultCollation.scala @@ -207,7 +207,7 @@ object ApplyDefaultCollation extends Rule[LogicalPlan] { newCreateView case createUserDefinedFunction@CreateUserDefinedFunction(ResolvedIdentifier( - catalog: SupportsNamespaces, identifier), _, _, _, _, _, _, _, _, _, _, _, _) + catalog: SupportsNamespaces, identifier), _, _, _, _, _, _, _, _, _, _, _, _, _) if createUserDefinedFunction.collation.isEmpty => val newCreateUserDefinedFunction = CurrentOrigin.withOrigin(createUserDefinedFunction.origin) { diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala index 45001dad6b08c..011d0f9f7caa4 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveCatalogs.scala @@ -92,7 +92,7 @@ class ResolveCatalogs(val catalogManager: CatalogManager) "CREATE", nameParts.last) case c @ CreateUserDefinedFunction( - u @ UnresolvedIdentifier(nameParts, _), _, _, _, _, _, _, _, _, _, _, _, _) => + u @ UnresolvedIdentifier(nameParts, _), _, _, _, _, _, _, _, _, _, _, _, _, _) => if (isSystemBuiltinName(nameParts)) { throw QueryCompilationErrors.operationNotAllowedOnBuiltinFunctionError( "CREATE", nameParts.last) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala index 32ab97c94ecf1..b15a99fe544b8 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala @@ -1644,6 +1644,7 @@ case class CreateUserDefinedFunction( containsSQL: Option[Boolean], language: org.apache.spark.sql.catalyst.catalog.RoutineLanguage, isTableFunc: Boolean, + isTemp: Boolean, ignoreIfExists: Boolean, replace: Boolean) extends UnaryCommand { override protected def withNewChildInternal(newChild: LogicalPlan): CreateUserDefinedFunction = diff --git a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala index fe52e4613b7e3..ce377be42d1f0 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala @@ -699,7 +699,7 @@ class ResolveSessionCatalog(val catalogManager: CatalogManager) throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) case c @ CreateUserDefinedFunction( - CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _) => + CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _, _) => CreateUserDefinedFunctionCommand( child = c.child, inputParamText = c.inputParamText, @@ -712,12 +712,12 @@ class ResolveSessionCatalog(val catalogManager: CatalogManager) containsSQL = c.containsSQL, language = c.language, isTableFunc = c.isTableFunc, - isTemp = false, + isTemp = c.isTemp, ignoreIfExists = c.ignoreIfExists, replace = c.replace) case CreateUserDefinedFunction( - ResolvedIdentifier(catalog, _), _, _, _, _, _, _, _, _, _, _, _, _) => + ResolvedIdentifier(catalog, _), _, _, _, _, _, _, _, _, _, _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala index 3a7e0bc5b3087..fcaba280d0bd1 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala @@ -1053,6 +1053,7 @@ class SparkSqlAstBuilder extends AstBuilder { containsSQL, language, isTableFunc, + isTemp = false, ctx.EXISTS != null, ctx.REPLACE != null) } else { @@ -1075,6 +1076,7 @@ class SparkSqlAstBuilder extends AstBuilder { containsSQL, language, isTableFunc, + isTemp = true, ctx.EXISTS != null, ctx.REPLACE != null) } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala index c6e7992e7d60d..3d57dc42417e8 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionParserSuite.scala @@ -61,6 +61,7 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL = containsSQL, language = org.apache.spark.sql.catalyst.catalog.LanguageSQL, isTableFunc = isTableFunc, + isTemp = false, ignoreIfExists = ignoreIfExists, replace = replace) } @@ -91,6 +92,7 @@ class CreateSQLFunctionParserSuite extends AnalysisTest { containsSQL = containsSQL, language = org.apache.spark.sql.catalyst.catalog.LanguageSQL, isTableFunc = isTableFunc, + isTemp = true, ignoreIfExists = ignoreIfExists, replace = replace) } From 95e3e4bee467e23d3352bad4ba239d4061c7631d Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Fri, 28 Aug 2026 14:57:49 +0530 Subject: [PATCH 11/14] [SPARK-50698][SQL] Enforce database=None for temporary function names in CreateSQLFunctionCommand --- .../command/CreateSQLFunctionCommand.scala | 23 ++++++++++++------- 1 file changed, 15 insertions(+), 8 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala index c424c2f6f9bf0..96fd5bfefc28a 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala @@ -68,14 +68,21 @@ case class CreateSQLFunctionCommand( import SQLFunction._ - lazy val name: FunctionIdentifier = child match { - case ResolvedIdentifier(c, ident) => - FunctionIdentifier(ident.name(), ident.namespace().headOption) - case u: UnresolvedIdentifier => - FunctionIdentifier(u.nameParts.last, u.nameParts.dropRight(1).lastOption) - case _ => - throw SparkException.internalError( - s"Unexpected child plan in CreateSQLFunctionCommand: $child") + lazy val name: FunctionIdentifier = { + val rawIdent = child match { + case ResolvedIdentifier(c, ident) => + FunctionIdentifier(ident.name(), ident.namespace().headOption) + case u: UnresolvedIdentifier => + FunctionIdentifier(u.nameParts.last, u.nameParts.dropRight(1).lastOption) + case _ => + throw SparkException.internalError( + s"Unexpected child plan in CreateSQLFunctionCommand: $child") + } + if (isTemp) { + FunctionIdentifier(rawIdent.funcName, None) + } else { + rawIdent + } } override protected def withNewChildInternal( From 4df08929586c76b4b50be29e814a36959625152c Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Sat, 29 Aug 2026 01:39:11 +0530 Subject: [PATCH 12/14] [SPARK-50698][SQL] Allow CreateUserDefinedFunction to create temporary functions regardless of catalog --- .../sql/catalyst/analysis/ResolveSessionCatalog.scala | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala index ce377be42d1f0..85a10609a5e5e 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/catalyst/analysis/ResolveSessionCatalog.scala @@ -698,8 +698,11 @@ class ResolveSessionCatalog(val catalogManager: CatalogManager) case CreateFunction(ResolvedIdentifier(catalog, _), _, _, _, _) => throw QueryCompilationErrors.missingCatalogCreateFunctionAbilityError(catalog) - case c @ CreateUserDefinedFunction( - CreateFunctionInSessionCatalog(ident), _, _, _, _, _, _, _, _, _, _, _, _, _) => + case c @ CreateUserDefinedFunction(child, _, _, _, _, _, _, _, _, _, _, _, _, _) + if c.isTemp || (child match { + case CreateFunctionInSessionCatalog(_) => true + case _ => false + }) => CreateUserDefinedFunctionCommand( child = c.child, inputParamText = c.inputParamText, From b6c40fd2e3539d85748766362925e23c549a6ef1 Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Sat, 29 Aug 2026 08:39:13 +0530 Subject: [PATCH 13/14] [SPARK-50698][SQL] Fix cyclic function reference check for temp functions and add MiMa filters --- project/MimaExcludes.scala | 22 +++++++++++++++---- .../command/CreateSQLFunctionCommand.scala | 12 ++++++++-- 2 files changed, 28 insertions(+), 6 deletions(-) diff --git a/project/MimaExcludes.scala b/project/MimaExcludes.scala index 6ad9255dc1bab..cd8b6942c4de2 100644 --- a/project/MimaExcludes.scala +++ b/project/MimaExcludes.scala @@ -72,10 +72,24 @@ object MimaExcludes { ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.status.api.v1.sql.Node.copy"), ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.status.api.v1.sql.Node$"), // [SPARK-50698][SQL] Refactor CreateUserDefinedFunctionCommand to extend from UnaryRunnableCommand - ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand"), - ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand.apply"), - ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.sql.execution.command.CreateSQLFunctionCommand"), - ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.apply"), + ProblemFilters.exclude[MissingTypesProblem]( + "org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand"), + ProblemFilters.exclude[DirectMissingMethodProblem]( + "org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand.apply"), + ProblemFilters.exclude[IncompatibleMethTypeProblem]( + "org.apache.spark.sql.execution.command.CreateUserDefinedFunctionCommand.apply"), + ProblemFilters.exclude[MissingTypesProblem]( + "org.apache.spark.sql.execution.command.CreateSQLFunctionCommand"), + ProblemFilters.exclude[DirectMissingMethodProblem]( + "org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.apply"), + ProblemFilters.exclude[IncompatibleMethTypeProblem]( + "org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.apply"), + ProblemFilters.exclude[IncompatibleMethTypeProblem]( + "org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.copy"), + ProblemFilters.exclude[IncompatibleMethTypeProblem]( + "org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.this"), + ProblemFilters.exclude[IncompatibleResultTypeProblem]( + "org.apache.spark.sql.execution.command.CreateSQLFunctionCommand.copy$default$1"), ProblemFilters.exclude[MissingClassProblem]( "org.apache.spark.sql.catalyst.plans.logical.CreateUserDefinedFunction"), ProblemFilters.exclude[MissingClassProblem]( diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala index 96fd5bfefc28a..02524d59ee7e2 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala @@ -433,7 +433,7 @@ case class CreateSQLFunctionCommand( } // Check cyclic reference using qualified function names. val newPath = path :+ f.function.name - if (f.function.name == name) { + if (isSameFunction(f.function.name)) { throw UserDefinedFunctionErrors.cyclicFunctionReference(newPath.mkString(" -> ")) } val plan = catalog.makeSQLTableFunctionPlan(f.name, f.function, f.inputs, f.output) @@ -444,6 +444,14 @@ case class CreateSQLFunctionCommand( } } + def isSameFunction(fName: FunctionIdentifier): Boolean = { + if (isTemp) { + fName.funcName == name.funcName + } else { + fName == name + } + } + def checkExpression(expression: Expression, path: Seq[FunctionIdentifier]): Unit = { expression.foreach { case s: SubqueryExpression => checkPlan(s.plan, path) @@ -456,7 +464,7 @@ case class CreateSQLFunctionCommand( } // Check cyclic reference using qualified function names. val newPath = path :+ f.function.name - if (f.function.name == name) { + if (isSameFunction(f.function.name)) { throw UserDefinedFunctionErrors.cyclicFunctionReference(newPath.mkString(" -> ")) } val plan = catalog.makeSQLFunctionPlan(f.name, f.function, f.inputs) From f6efd8d2f36ce0db36f553e4efece8ea3122148c Mon Sep 17 00:00:00 2001 From: zahed1994 Date: Sat, 29 Aug 2026 15:34:53 +0530 Subject: [PATCH 14/14] [SPARK-50698][SQL][FOLLOWUP] Include catalog in FunctionIdentifier extraction for ResolvedIdentifier --- .../command/CreateSQLFunctionCommand.scala | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala index 02524d59ee7e2..602e4966168a3 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/CreateSQLFunctionCommand.scala @@ -71,15 +71,22 @@ case class CreateSQLFunctionCommand( lazy val name: FunctionIdentifier = { val rawIdent = child match { case ResolvedIdentifier(c, ident) => - FunctionIdentifier(ident.name(), ident.namespace().headOption) + FunctionIdentifier(ident.name(), ident.namespace().headOption, Some(c.name())) case u: UnresolvedIdentifier => - FunctionIdentifier(u.nameParts.last, u.nameParts.dropRight(1).lastOption) + val parts = u.nameParts + if (parts.length >= 3) { + FunctionIdentifier(parts.last, Some(parts(parts.length - 2)), Some(parts.head)) + } else if (parts.length == 2) { + FunctionIdentifier(parts.last, Some(parts.head), None) + } else { + FunctionIdentifier(parts.last, None, None) + } case _ => throw SparkException.internalError( s"Unexpected child plan in CreateSQLFunctionCommand: $child") } if (isTemp) { - FunctionIdentifier(rawIdent.funcName, None) + FunctionIdentifier(rawIdent.funcName, None, None) } else { rawIdent }