Apache Spark SQL 的 CREATE FUNCTION (SQL) 完整指南:标量函数、表函数与 SQL UDF 最佳实践
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
本文基于 Apache Spark 官方 SQL 参考文档编写,全面讲解CREATE FUNCTION (SQL)语句——用于在 Spark SQL 中定义以纯 SQL 表达式或查询为函数体的用户自定义函数(SQL UDF)。读者将掌握临时函数与永久函数的作用域差异、标量函数与表函数(Table Function)的定义方式、函数参数与特征(characteristic)子句的完整语法,以及 SQL Path 冻结机制、函数替换与元数据检查等进阶用法,并了解该语句在 Spark 源码中的解析与执行链路。
概述:什么是 SQL 函数
CREATE FUNCTION语句用于创建一个可以在 SQL 语句中直接调用的函数。它有两种形态:
- 临时(TEMPORARY)函数:仅在当前会话内可见,会话结束即被丢弃,不会在 catalog 中留下持久化条目,位于会话级
system.session命名空间中; - 永久(Permanent)函数:持久化到 catalog 中,跨会话可用。
函数的返回值可以是标量(scalar),也可以是表(table);函数体可以是一个 SQL 表达式(expression),也可以是一个查询(query)。OR REPLACE选项允许更新已存在的函数定义,IF NOT EXISTS则避免重复创建时报错。
在 Spark 源码的 ANTLR 语法定义中(SqlBaseParser.g4),CREATE FUNCTION被区分为两条规则:
#createFunction:形如CREATE FUNCTION name AS class (USING jar),用于注册 JVM 外部实现函数;#createUserDefinedFunction:形如CREATE FUNCTION name(params) RETURNS ... RETURN expression/query,即本文主题的SQL 函数。
对应的逻辑计划节点为 CreateUserDefinedFunction,它携带输入参数文本、返回类型文本、表达式/查询文本、确定性标记、是否包含 SQL 数据、函数语言、是否为表函数、IF NOT EXISTS与replace标记等字段,是 SQL 函数从语法到执行的核心中间表示。
语法结构
CREATE [OR REPLACE] [TEMPORARY] FUNCTION [IF NOT EXISTS] function_name ( [ function_parameter [, ...] ] ) { [ RETURNS data_type ] | RETURNS TABLE [ ( column_spec [, ...]) ] } [ characteristic [...] ] RETURN { expression | query } function_parameter parameter_name data_type [DEFAULT default_expression] [COMMENT parameter_comment] column_spec column_name data_type [COMMENT column_comment] characteristic { LANGUAGE SQL | [NOT] DETERMINISTIC | COMMENT function_comment | [CONTAINS SQL | READS SQL DATA] }从语法可见,函数定义包含四大部分:可选修饰子句(OR REPLACE/TEMPORARY/IF NOT EXISTS)、函数名与参数列表、返回类型声明、函数体(RETURN关键字后跟表达式或查询)。
参数详解
OR REPLACE
指定后,将替换同名且同签名(参数个数与参数类型相同)的函数。该选项主要用于更新函数体和返回类型。
需要注意:
- 不能替换签名不同的函数,也不能替换存储过程(procedure);
OR REPLACE与IF NOT EXISTS互斥,不能同时指定。
TEMPORARY
控制函数作用域:
- 指定时,函数仅在当前会话有效,位于会话级
system.session命名空间; - 会话结束即被丢弃,catalog 中不产生持久化条目;
- 未指定时,函数持久化到 catalog,跨会话可用。
IF NOT EXISTS
- 仅当函数不存在时才创建;
- 若同名函数已存在,创建操作静默成功(不抛错);
- 与
OR REPLACE互斥。
function_name
函数命名规则因函数类型而异:
- 永久函数:可用数据库名(或 catalog + 数据库)限定,语法为
[ catalog_name. ] [ database_name. ] function_name;若不加限定,则创建在当前 schema 中。 - 临时函数:可用
session或system.session限定,语法为[ { session | system.session } . ] function_name;任何其他限定符(包括system.builtin、当前 schema、任意数据库名)都会被拒绝,抛出INVALID_TEMP_OBJ_QUALIFIER错误。
函数名在其所在 schema 的所有例程(procedure 与 function)中必须唯一。
function_parameter
定义函数的一个输入参数:
- parameter_name:参数名,函数内必须唯一,可参考标识符文档;
- data_type:任意受支持的数据类型;
- DEFAULT default_expression:可选默认值,当调用方未给该参数传值时使用。
default_expression必须能 cast 为data_type,且不能引用其他参数或包含子查询;一旦某个参数指定了默认值,其后的所有参数也必须指定默认值; - COMMENT parameter_comment:可选的参数说明,必须是
STRING字面量。
RETURNS data_type
标量函数的返回数据类型。该子句可选:若省略,则从 SQL 函数体推导返回类型。
RETURNS TABLE [ ( column_spec [, ...] ) ]
标记该函数为表函数,并可选地声明结果表的签名:
- column_name:列名,签名内必须唯一;
- data_type:任意受支持的数据类型;
- COMMENT column_comment:可选的列说明,必须是
STRING字面量。
若未指定column_spec,结果签名将从 SQL UDF 函数体推导。
RETURN { expression | query }
函数体:
- 标量函数:函数体可以是表达式或查询;
- 表函数:函数体只能是查询。
表达式内部不允许包含:
- 聚合函数;
- 窗口函数;
- 排名函数;
- 产生多行的函数(如
explode)。
其他约束与语义:
- 永久 SQL UDF不能引用临时视图、临时函数或会话变量;
- SQL Path 冻结:
CREATE FUNCTION时刻生效的 SQL Path 会被捕获进函数元数据,之后每次调用都按该冻结 Path 解析函数体,而不是调用方当前会话的 Path。但函数体内的current_schema()与current_path()仍返回调用方的上下文。可用 DESCRIBE FUNCTION EXTENDED 查看捕获的 Path,Path 的设置见 SET PATH; - 函数体内引用参数时,可用参数的无限定名,也可用“函数名.参数名”的方式限定(例如表函数示例中的
weekdays.start)。
characteristic(特征子句)
所有特征子句均可选,可任意数量、任意顺序出现,但每个子句只能出现一次:
- LANGUAGE SQL:声明函数实现语言为 SQL;
- [NOT] DETERMINISTIC:声明函数是否确定。确定性函数指给定一组参数只返回一个结果的函数。你可以主动标注(例如函数体非确定但标记为
DETERMINISTIC以鼓励常量折叠等查询优化,或反向抑制缓存);未指定时从函数体推导; - COMMENT function_comment:函数注释,必须是
STRING字面量; - CONTAINS SQL / READS SQL DATA:声明函数是否直接或间接读取表/视图的数据。若函数读取 SQL 数据则不能指定
CONTAINS SQL;两者都未指定时从函数体推导。
完整示例:从标量函数到表函数
以下示例均来自官方文档,可直接在 Spark SQL 中运行验证。
创建并使用 SQL 标量函数
> CREATE VIEW t(c1, c2) AS VALUES (0, 1), (1, 2); -- 创建无参数的临时函数 > CREATE TEMPORARY FUNCTION hello() RETURNS STRING RETURN 'Hello World!'; > SELECT hello(); Hello World! -- 创建带参数的永久函数 > CREATE FUNCTION area(x DOUBLE, y DOUBLE) RETURNS DOUBLE RETURN x * y; -- 在 SELECT 子句中使用 SQL 函数 > SELECT area(c1, c2) AS area FROM t; 1.0 1.0 -- 在 WHERE 子句中使用 SQL 函数 > SELECT * FROM t WHERE area(c1, c2) > 0; 1 2 -- 组合 SQL 函数(函数体中调用其他函数) > CREATE FUNCTION square(x DOUBLE) RETURNS DOUBLE RETURN area(x, x); > SELECT c1, square(c1) AS square FROM t; 0 0.0 1 1.0 -- 创建非确定函数(模拟掷骰子) > CREATE FUNCTION roll_dice() RETURNS INT NOT DETERMINISTIC CONTAINS SQL COMMENT 'Roll a single 6 sided die' RETURN (rand() * 6)::INT + 1; -- 掷一个六面骰子 > SELECT roll_dice(); 3roll_dice示例完整演示了特征子句的用法:NOT DETERMINISTIC告知优化器该函数每次调用结果可能不同(从而抑制常量折叠),CONTAINS SQL说明函数体不直接读取表数据,COMMENT为函数附加说明。
创建 SQL 表函数
表函数的函数体必须是查询,且可以使用LATERAL VIEW、sequence等复杂构造:
-- 生成两个日期之间的所有工作日 > CREATE FUNCTION weekdays(start DATE, end DATE) RETURNS TABLE(day_of_week STRING, day DATE) RETURN SELECT extract(DAYOFWEEK_ISO FROM day), day FROM (SELECT sequence(weekdays.start, weekdays.end)) AS T(days) LATERAL VIEW explode(days) AS day WHERE extract(DAYOFWEEK_ISO FROM day) BETWEEN 1 AND 5; -- 返回所有工作日 > SELECT weekdays.day_of_week, day FROM weekdays(DATE'2022-01-01', DATE'2022-01-14'); 1 2022-01-03 2 2022-01-04 3 2022-01-05 4 2022-01-06 5 2022-01-07 1 2022-01-10 2 2022-01-11 3 2022-01-12 4 2022-01-13 5 2022-01-14 -- 从 LATERAL 关联的日期区间调用表函数 > SELECT weekdays.* FROM VALUES (DATE'2020-01-01'), (DATE'2021-01-01'), (DATE'2022-01-01') AS starts(start), LATERAL weekdays(start, start + INTERVAL '7' DAYS); 3 2020-01-01 4 2020-01-02 5 2020-01-03 1 2020-01-06 2 2020-01-07 3 2020-01-08 5 2021-01-01 1 2021-01-04 2 2021-01-05 3 2021-01-06 4 2021-01-07 5 2021-01-08 1 2022-01-03 2 2022-01-04 3 2022-01-05 4 2022-01-06 5 2022-01-07注意表函数中weekdays.start/weekdays.end这种“函数名限定参数”的写法,以及LATERAL weekdays(...)的调用方式——表函数可以与LATERAL关联(LATERAL correlation)结合,对每一行输入动态展开为多行输出,这正是表函数区别于标量函数的核心价值。
替换 SQL 函数
-- 替换 SQL 标量函数 > CREATE OR REPLACE FUNCTION square(x DOUBLE) RETURNS DOUBLE RETURN x * x; -- 替换 SQL 表函数 > CREATE OR REPLACE FUNCTION getemps(deptno INT) RETURNS TABLE (name STRING) RETURN SELECT name FROM employee e WHERE e.deptno = getemps.deptno; -- 描述一个 SQL 表函数 > DESCRIBE FUNCTION getemps; Function: default.getemps Type: TABLE Input: deptno INT Returns: id INT name STRINGOR REPLACE要求新旧函数签名一致,适用于升级函数体逻辑或修正返回类型。
描述 SQL 函数
使用 DESCRIBE FUNCTION 查看函数元数据(类型、输入、返回):
> DESCRIBE FUNCTION hello; Function: hello Type: SCALAR Input: () Returns: STRING > DESCRIBE FUNCTION area; Function: default.area Type: SCALAR Input: x DOUBLE y DOUBLE Returns: DOUBLE > DESCRIBE FUNCTION roll_dice; Function: default.roll_dice Type: SCALAR Input: num_dice INT num_sides INT Returns: INT创建带会话限定符的临时 SQL 函数
临时函数可以使用session/system.session限定,三种写法指向同一个函数:
-- 不加限定、加 session 限定、加 system.session 限定,都创建同一个临时函数 > CREATE TEMPORARY FUNCTION add_one(x INT) RETURNS INT RETURN x + 1; > CREATE OR REPLACE TEMPORARY FUNCTION session.add_one(x INT) RETURNS INT RETURN x + 1; > CREATE OR REPLACE TEMPORARY FUNCTION system.session.add_one(x INT) RETURNS INT RETURN x + 1; -- 三个名字指向同一个临时函数: > SELECT add_one(1), session.add_one(1), system.session.add_one(1); 2 2 2 -- DROP TEMPORARY FUNCTION 也接受同样的限定符: > DROP TEMPORARY FUNCTION session.add_one; -- 任何其他限定符都会被拒绝: > CREATE TEMPORARY FUNCTION mydb.bad_temp() RETURNS INT RETURN 1; [INVALID_TEMP_OBJ_QUALIFIER] qualifier `mydb` is not allowed for temporary FUNCTION ... > CREATE TEMPORARY FUNCTION system.builtin.bad_temp() RETURNS INT RETURN 1; [INVALID_TEMP_OBJ_QUALIFIER] qualifier `system`.`builtin` is not allowed for temporary FUNCTION ...Frozen SQL Path(冻结 SQL Path)机制
SQL UDF 在CREATE FUNCTION时刻捕获当前生效的 SQL Path,之后每次调用都按冻结的 Path 解析函数体中的对象引用,即使调用方会话后来设置了不同的 PATH:
> CREATE SCHEMA path_a; > CREATE SCHEMA path_b; > CREATE TABLE path_a.t USING parquet AS SELECT 10 AS id; > CREATE TABLE path_b.t USING parquet AS SELECT 20 AS id; -- 创建函数时 PATH 指向 path_a,所以函数体里无限定的 `t` 绑定到 path_a.t > SET PATH = spark_catalog.path_a, system.builtin; > CREATE FUNCTION default.frozen_fn() RETURNS INT RETURN (SELECT MAX(id) FROM t); -- 切换当前 PATH > SET PATH = spark_catalog.path_b, system.builtin; -- 裸查询跟随“实时”PATH: > SELECT MAX(id) FROM t; 20 -- 函数体跟随其“冻结”的 PATH: > SELECT default.frozen_fn(); 10 -- DESCRIBE FUNCTION EXTENDED 展示捕获的 PATH: > DESC FUNCTION EXTENDED default.frozen_fn; Function: spark_catalog.default.frozen_fn ... SQL Path: spark_catalog.path_a, system.builtin这是 SQL 函数可复现性的关键设计:函数行为不随调用方会话的 PATH 变化而漂移,保证函数在任何会话中返回一致的结果。相关的 Path 管理语法详见 SET PATH,名称解析规则可参考Name Resolution。
源码视角:SQL 函数的解析与执行链路
从当前仓库源码可以完整还原 SQL 函数的处理链路:
语法层:在 SqlBaseParser.g4 中,
CREATE ... FUNCTION语法通过#createUserDefinedFunction标签与#createFunction(外部 JVM 函数)区分开来,SQL 函数的解析入口要求参数列表LEFT_PAREN ... RIGHT_PAREN、可选的RETURNS声明、routineCharacteristics(即特征子句)以及RETURN (query | expression)。AST 构建与逻辑计划:AstBuilder 将语法树转换为 CreateUserDefinedFunction 逻辑计划节点,其中
isDeterministic、containsSQL、language、isTableFunc、ignoreIfExists、replace等字段一一对应本文介绍的语法选项。校验与执行:解析器测试 CreateSQLFunctionParserSuite 覆盖了完整的语法面,包括:无参/带参函数、标量与表函数、
OR REPLACE、IF NOT EXISTS、临时函数的命名空间限定(a.b/a.b.c)、参数 COMMENT,以及非法用法(如OR REPLACE与IF NOT EXISTS混用、LANGUAGE blah、SPECIFIC、NO SQL、重复的特征子句等)的报错。端到端行为验证:仓库的 SQL 测试套件(如 sql-udf.sql、sql-path.sql 及其 results/analyzer-results 对应输出)通过 golden 文件验证 SQL UDF 的创建、调用、名字优先级与 PATH 解析行为。
如果需要在测试或实践中验证语法解析,可参考上述测试文件中的断言模式;语法解析器的完整单元测试位于 CreateSQLFunctionParserSuite。
常见约束与注意事项
综合语法、参数说明与源码测试,使用 SQL 函数时有以下要点值得注意:
OR REPLACE与IF NOT EXISTS互斥,同时指定会直接报语法错误;- 临时函数只能使用
session/system.session限定,system.builtin、普通库名等限定符均报INVALID_TEMP_OBJ_QUALIFIER; - 永久 SQL UDF 不能依赖临时对象(临时视图、临时函数、会话变量),否则函数在跨会话调用时会解析失败;
- 参数默认值有位置约束:一旦某参数设置
DEFAULT,其后所有参数都必须设置;默认表达式不能引用其他参数或子查询; - 表函数体必须是查询,标量函数体可以是表达式或查询;函数体禁止聚合函数、窗口/排名函数及行产生函数(如
explode)——注意表函数示例中的explode位于函数体查询内部,属于查询内合法使用; - 特征子句每类只能出现一次;
READS SQL DATA与CONTAINS SQL互斥; - SQL Path 在创建时冻结:若函数依赖的未限定对象随后被移动到其他 schema,函数仍解析到创建时的对象。
相关语句
SQL 函数的管理与探查与以下语句配套使用:
- SHOW FUNCTIONS:列出函数;
- DESCRIBE FUNCTION:查看函数签名与类型;
- DROP FUNCTION:删除函数;
- SET PATH:管理函数体解析所用的 SQL Path;
- Name Resolution:理解对象名(含函数名)的解析顺序与规则。
总结
CREATE FUNCTION (SQL)让用户无需编写 Java/Scala/Python 代码,仅用纯 SQL 即可定义可复用的标量函数与表函数:临时函数服务于会话内的快速封装,永久函数沉淀为跨会话共享的资产;DEFAULT参数、特征子句与冻结 SQL Path 等设计,为函数提供了接近一等公民的表达力与可复现性。结合 DESCRIBE FUNCTION 与 DROP FUNCTION 形成完整生命周期管理,是构建 Spark SQL 业务抽象层、报表计算层与数据治理层时值得优先考虑的内建方案。
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考