StarRocks Java UDF 开发实战:从标量函数到 UDAF/UDWF/UDTF 的完整指南
2026/9/17 10:14:46 网站建设 项目流程

StarRocks Java UDF 开发实战:从标量函数到 UDAF/UDWF/UDTF 的完整指南

【免费下载链接】starrocksThe world's fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks

Java UDF(User-Defined Function)是 StarRocks 提供给用户的扩展机制,允许使用 Java 语言编译自定义函数以满足特定业务需求。本文基于官方文档与仓库源码,系统讲解标量 UDF、用户自定义聚合函数(UDAF)、用户自定义窗口函数(UDWF)与用户自定义表函数(UDTF)的完整开发、打包、注册与使用流程,同时深入剖析底层实现原理与向量化(Arrow)输入等进阶特性。读完本文,你将能够独立创建一个 Maven 项目,编写四类 Java UDF 并在 StarRocks 集群中注册与调用。

Java UDF 功能概览

自 StarRocksv2.2.0起,用户可以编译 Java UDF 以满足特定业务需求;自v3.0起,StarRocks 支持全局 UDF(Global UDF),只需在相关 SQL 语句(CREATE/SHOW/DROP)中加入GLOBAL关键字即可。

目前 StarRocks 支持以下四类 Java UDF:

UDF 类型全称行为特征
Scalar UDF标量函数单行输入,单值输出
UDAF用户自定义聚合函数多行输入,单值输出
UDWF用户自定义窗口函数按窗口(OVER 子句划分的行集)逐行计算
UDTF用户自定义表函数单行输入,多行(表)输出,常用于行转列

从源码结构看,BE 端为四类函数分别维护了执行入口:java_function_call_expr.cpp(标量)、java_udaf_function.cpp(聚合)、java_window_function.cpp(窗口)与 java_udtf_function.cpp(表函数),它们统一基于 java_udf.cpp 的 JNI 调用层与嵌入 JVM 交互。

前置条件

在开始开发 Java UDF 之前,需要满足以下条件:

  1. 安装 Apache Maven,用于创建和编译 Java 工程。
  2. 服务器安装 JDK 17
  3. 开启 Java UDF 特性:在 FE 配置文件fe/conf/fe.conf中将 FE 配置项enable_udf设置为true,然后重启 FE 节点使其生效。详细配置说明参见 FE 参数配置。

关于该配置项,仓库源码中有明确佐证:Config.java 中声明public static boolean enable_udf = false;,即默认关闭,必须显式开启才能使用 Java UDF 特性。同时建议查看 udf_security.policy,FE 的 JAVA_OPTS 中通过-Djava.security.policy=${STARROCKS_HOME}/conf/udf_security.policy(见 conf/fe.conf)加载该 Java 安全策略文件,默认授予java.security.AllPermission,可在此基础上按需收紧 UDF 运行权限。

开发与使用 Java UDF 的完整流程

整体流程为:创建 Maven 工程 → 添加依赖 → 编写 Java 类 → 打包 → 上传 JAR → 在 StarRocks 中注册函数 → 在 SQL 中调用。

Step 1:创建 Maven 项目

创建一个 Maven 项目,其基本目录结构如下:

project |--pom.xml |--src | |--main | | |--java | | |--resources | |--test |--target

仓库中的 contrib/java-udf 目录就是一个可直接参考的完整示例工程,其pom.xml位于 contrib/java-udf/pom.xml,源码位于 contrib/java-udf/src/main/java/com/starrocks/udf。

Step 2:添加依赖

pom.xml文件中添加如下依赖。官方示例使用com.alibaba:fastjson(用于 JSON 解析类 UDF),并配置maven-dependency-plugin将依赖拷贝到target/lib,同时配置maven-assembly-plugin打一个包含所有依赖的 fat JAR:

<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>org.example</groupId> <artifactId>udf</artifactId> <version>1.0-SNAPSHOT</version> <properties> <maven.compiler.source>17</maven.compiler.source> <maven.compiler.target>17</maven.compiler.target> </properties> <dependencies> <dependency> <groupId>com.alibaba</groupId> <artifactId>fastjson</artifactId> <version>1.2.76</version> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-dependency-plugin</artifactId> <version>2.10</version> <executions> <execution> <id>copy-dependencies</id> <phase>package</phase> <goals> <goal>copy-dependencies</goal> </goals> <configuration> <outputDirectory>${project.build.directory}/lib</outputDirectory> </configuration> </execution> </executions> </plugin> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-assembly-plugin</artifactId> <version>3.3.0</version> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> <configuration> <descriptorRefs> <descriptorRef>jar-with-dependencies</descriptorRef> </descriptorRefs> </configuration> </plugin> </plugins> </build> </project>

注意maven.compiler.sourcemaven.compiler.target需设置为 17,与 JDK 17 前置条件保持一致。仓库示例工程 contrib/java-udf/pom.xml 同样使用java.version=17

Step 3:编写 UDF Java 类

使用 Java 语言编写 UDF。方法中的请求参数与返回参数的数据类型必须与 Step 6 中CREATE FUNCTION语句声明的类型一致,并遵循下文「SQL 数据类型与 Java 数据类型映射」中的对应关系。

编写标量 UDF(Scalar UDF)

标量 UDF 对单行数据进行操作并返回单个值,查询结果集中每一行对应一个值。典型的内置标量函数包括UPPERLOWERROUNDABS

业务场景示例:假设 JSON 数据中某个字段的值是 JSON 字符串而非 JSON 对象。使用 SQL 提取 JSON 字符串时需要嵌套执行两次GET_JSON_STRING,例如GET_JSON_STRING(GET_JSON_STRING('{"key":"{\"k0\":\"v0\"}"}', "$.key"), "$.k0")。为简化 SQL,可以编写一个能直接提取 JSON 字符串的标量 UDF,例如MY_UDF_JSON_GET('{"key":"{\"k0\":\"v0\"}"}', "$.key.k0")

package com.starrocks.udf.sample; import com.alibaba.fastjson.JSONPath; public class UDFJsonGet { public final String evaluate(String obj, String key) { if (obj == null || key == null) return null; try { // JSONPath 库可以完整展开字段值中的 JSON 字符串 return JSONPath.read(obj, key).toString(); } catch (Exception e) { return null; } } }

用户自定义类必须实现下表中的方法:

方法说明
TYPE1 evaluate(TYPE2, ...)执行 UDF。evaluate()方法需要public访问级别
编写 UDAF(用户自定义聚合函数)

UDAF 对多行数据进行操作并返回单个值。典型的内置聚合函数包括SUMCOUNTMAXMIN,它们聚合每个GROUP BY子句指定的多行数据并返回单个值。

业务场景示例:编写名为MY_SUM_INT的 UDAF。与返回 BIGINT 类型值的内置聚合函数SUM不同,MY_SUM_INT只支持 INT 数据类型的请求参数和返回参数。

package com.starrocks.udf.sample; public class SumInt { public static class State { int counter = 0; public int serializeLength() { return 4; } } public State create() { return new State(); } public void destroy(State state) { } public final void update(State state, Integer val) { if (val != null) { state.counter+= val; } } public void serialize(State state, java.nio.ByteBuffer buff) { buff.putInt(state.counter); } public void merge(State state, java.nio.ByteBuffer buffer) { int val = buffer.getInt(); state.counter += val; } public Integer finalize(State state) { return state.counter; } }

用户自定义类必须实现下表中的方法:

方法说明
State create()创建一个状态(State)
void destroy(State)销毁一个状态
void update(State, ...)更新状态。除第一个参数State外,还可以在 UDF 声明中指定一个或多个请求参数
void serialize(State, ByteBuffer)将状态序列化到字节缓冲区
void merge(State, ByteBuffer)从字节缓冲区反序列化出一个状态,并合并到第一个参数指定的状态中
TYPE finalize(State)从状态中获取 UDF 的最终结果

编译 UDAF 时,还必须使用缓冲区类java.nio.ByteBuffer和局部变量serializeLength

类与局部变量说明
java.nio.ByteBuffer()缓冲区类,用于存储中间结果。中间结果在节点间传输执行时可能被序列化或反序列化,因此还必须使用serializeLength变量指定中间结果反序列化时允许的长度
serializeLength()中间结果反序列化时允许的长度,单位:字节。该局部变量应设置为 INT 类型值。例如State { int counter = 0; public int serializeLength() { return 4; }}表示中间结果为 INT 类型、反序列化长度为 4 字节。可根据业务需求调整:例如希望中间结果为 LONG 类型、反序列化长度为 8 字节,则传入State { long counter = 0; public int serializeLength() { return 8; }}

关于存储在java.nio.ByteBuffer类中的中间结果的反序列化,需要注意以下几点:

  • 不能调用ByteBuffer类依赖的remaining()方法来反序列化状态。
  • 不能对ByteBuffer类调用clear()方法。
  • serializeLength的值必须与写入数据的长度一致,否则序列化与反序列化会产生错误结果。

源码佐证:仓库 contrib/java-udf/src/main/java/com/starrocks/udf/SumMap.java 是一个更复杂的 Map 求和 UDAF 示例,其State类同时实现了serialize()/serializeLength()/deserialize()等辅助方法,展示了自定义序列化逻辑的写法;SumMapInt64.java 则是其 Long 值变体。

编写 UDWF(用户自定义窗口函数)

与常规聚合函数不同,UDWF 对一组多行数据(统称为窗口)进行操作,并为每一行返回一个值。典型窗口函数通过OVER子句将行划分为多个集合,对每个集合执行计算,并为每一行返回一个值。

业务场景示例:编写名为MY_WINDOW_SUM_INT的 UDWF。与返回 BIGINT 类型值的内置聚合函数SUM不同,MY_WINDOW_SUM_INT只支持 INT 数据类型的请求参数和返回参数。

package com.starrocks.udf.sample; public class WindowSumInt { public static class State { int counter = 0; public int serializeLength() { return 4; } @Override public String toString() { return "State{" + "counter=" + counter + '}'; } } public State create() { return new State(); } public void destroy(State state) { } public void update(State state, Integer val) { if (val != null) { state.counter+=val; } } public void serialize(State state, java.nio.ByteBuffer buff) { buff.putInt(state.counter); } public void merge(State state, java.nio.ByteBuffer buffer) { int val = buffer.getInt(); state.counter += val; } public Integer finalize(State state) { return state.counter; } public void reset(State state) { state.counter = 0; } public void windowUpdate(State state, int peer_group_start, int peer_group_end, int frame_start, int frame_end, Integer[] inputs) { for (int i = (int)frame_start; i < (int)frame_end; ++i) { state.counter += inputs[i]; } } }

用户自定义类必须实现 UDAF 要求的所有方法(因为 UDWF 是一种特殊的聚合函数),以及下表中的windowUpdate()方法:

方法说明
void windowUpdate(State state, int, int, int, int, ...)更新窗口数据。有关 UDWF 的更多信息,参见 Window functions。每次输入一行数据时,该方法都会获取窗口信息并相应更新中间结果。
  • peer_group_start:当前分区的起始位置。OVER子句使用PARTITION BY指定分区列,分区列值相同的行视为同一分区。
  • peer_group_end:当前分区的结束位置。
  • frame_start:当前窗口帧的起始位置。窗口帧子句指定计算范围,覆盖当前行及与当前行指定距离内的行。例如ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING指定计算范围覆盖当前行、当前行的前一行和当前行的后一行。
  • frame_end:当前窗口帧的结束位置。
  • inputs:输入窗口的数据。数据是一个数组包,只支持特定数据类型。本示例中输入为 INT 值,数组包为Integer[]
编写 UDTF(用户自定义表函数)

UDTF 读取一行数据并返回多个值,这些值可以被视为一张表。表函数通常用于将行转换为列。

注意:StarRocks 允许 UDTF 返回一张由多行和一列组成的表。

业务场景示例:编写名为MY_UDF_SPLIT的 UDTF。MY_UDF_SPLIT以空格为分隔符,请求参数和返回参数均为 STRING 数据类型。

package com.starrocks.udf.sample; public class UDFSplit{ public String[] process(String in) { if (in == null) return null; return in.split(" "); } }

用户自定义类定义的方法必须满足以下要求:

方法说明
TYPE[] process()执行 UDTF 并返回一个数组

Step 4:打包 Java 工程

运行以下命令打包:

mvn package

target目录下会生成两个 JAR 文件:udf-1.0-SNAPSHOT.jarudf-1.0-SNAPSHOT-jar-with-dependencies.jar

Step 5:上传 Java 工程

udf-1.0-SNAPSHOT-jar-with-dependencies.jar(fat JAR)上传到一个持续运行、且集群内所有 FE 和 BE 均可访问的 HTTP 服务器上,然后执行以下命令部署该文件:

mvn deploy

也可以使用 Python 搭建一个简单的 HTTP 服务器来托管 JAR 文件。

注意:在 Step 6 中,FE 会检查包含 UDF 代码的 JAR 文件并计算校验和,BE 会下载并执行该 JAR 文件。因此 HTTP 服务器必须保持在线,且所有 FE、BE 节点都能访问到该 URL。

Step 6:在 StarRocks 中创建 UDF

StarRocks 支持在两种命名空间下创建 UDF:数据库命名空间全局命名空间

  • 如果 UDF 没有可见性或隔离性需求,可以创建为全局 UDF。之后可以直接使用函数名引用,无需在函数名前添加 catalog 和数据库名前缀。
  • 如果 UDF 有可见性或隔离性需求,或者需要在不同数据库中创建相同的 UDF,可以在每个数据库下分别创建。会话连接到目标数据库时,可以直接使用函数名引用;会话连接到其他 catalog 或数据库时,需要以catalog.database.function的形式加上前缀引用。

NOTICE:创建和使用全局 UDF 之前,必须联系系统管理员授予所需权限,参见 GRANT。

上传 JAR 包后,即可在 StarRocks 中创建 UDF。对于全局 UDF,创建语句中必须包含GLOBAL关键字。

创建语法
CREATE [GLOBAL][AGGREGATE | TABLE] FUNCTION function_name (arg_type [, ...]) RETURNS return_type PROPERTIES ("key" = "value" [, ...])
参数说明
参数是否必选说明
GLOBAL是否创建全局 UDF,v3.0 起支持
AGGREGATE是否创建 UDAF 或 UDWF
TABLE是否创建 UDTF。若AGGREGATETABLE都未指定,则创建标量函数
function_name要创建的函数名,可包含数据库名,例如db1.my_func。若function_name包含数据库名,则在对应数据库创建 UDF,否则在当前数据库创建。新函数名与其参数不能与目标数据库中已有名称相同,否则创建失败;函数名相同但参数不同时,创建成功
arg_type函数参数类型,可用, ...表示多个参数。支持的数据类型见「SQL 数据类型与 Java 数据类型映射」
return_type函数返回类型,支持的数据类型同上
PROPERTIES函数属性,根据创建的 UDF 类型不同而不同
创建标量 UDF

执行以下命令创建上文编译的标量 UDF:

CREATE [GLOBAL] FUNCTION MY_UDF_JSON_GET(string, string) RETURNS string PROPERTIES ( "symbol" = "com.starrocks.udf.sample.UDFJsonGet", "type" = "StarrocksJar", "file" = "http://http_host:http_port/udf-1.0-SNAPSHOT-jar-with-dependencies.jar" );

PROPERTIES 参数说明:

参数说明
symbolUDF 所属 Maven 工程的类名,格式为<package_name>.<class_name>
typeUDF 类型。设置为StarrocksJar,表示该 UDF 是基于 Java 的函数
file下载包含 UDF 代码的 JAR 文件的 HTTP URL,格式为http://<http_server_ip>:<http_server_port>/<jar_package_name>
isolation可选。若要在多次 UDF 执行之间共享函数实例并支持静态变量,则设置为"shared"
input可选。输入格式,合法值为scalar(默认,每行一个装箱的 Java 对象)与arrow(向量化,每个参数一个 Apache ArrowFieldVector,覆盖整批数据)。参见下文「向量化(Arrow)输入」

源码佐证:FE 端 CreateFunctionAnalyzer.java 负责解析CREATE FUNCTION语句:读取symbolinputisolation等属性(如CreateFunctionStmt.SYMBOL_KEYISOLATION_KEY),校验input只允许"arrow""scalar",并分别调用analyzeStarrocksJarUdf/analyzeStarrocksJarUdaf/analyzeStarrocksJarUdtf对三类函数做类型与类结构校验。创建 UDF 时 FE 还会下载 JAR 并计算校验和,BE 执行时再次校验,确保函数代码未被篡改。

向量化(Arrow)输入

"input"设置为"arrow"可以把 Java UDF 从「逐行装箱」的调用约定切换为向量化调用约定:你的方法将为每个参数接收一个覆盖整批数据的 Apache ArrowFieldVector,标量/UDTF 则返回一个FieldVector。列数据通过 Arrow C Data Interface 与后端零拷贝交换,避免了默认路径下逐行的装箱/拆箱开销。

首先在 Maven 工程中添加 Arrow 依赖(版本必须与你 StarRocks 发行版内置的版本一致),scope 设为provided(BE 已自带该依赖):

<dependency> <groupId>org.apache.arrow</groupId> <artifactId>arrow-vector</artifactId> <version>17.0.0</version> <scope>provided</scope> </dependency>

结果向量请通过arg.getAllocator()从框架托管的 allocator 中分配,这样引擎可以接管其生命周期。

注意:BE 在嵌入式 JVM 中运行 UDF,该 JVM 必须向 Arrow 的堆外内存层暴露java.nio。仓库自带的 conf/be.conf 已通过JAVA_OPTS="... --add-opens=java.base/java.nio=ALL-UNNAMED ..."完成该设置;如果你自定义JAVA_OPTS,必须保留该 flag,否则 arrow-input UDF 将初始化失败。

标量 UDF——evaluate(FieldVector...)返回一个FieldVector

public class ArrowAdd { public IntVector evaluate(IntVector a, IntVector b) { IntVector out = new IntVector("result", a.getAllocator()); int n = a.getValueCount(); out.allocateNew(n); for (int i = 0; i < n; i++) { if (a.isNull(i) || b.isNull(i)) { out.setNull(i); } else { out.set(i, a.get(i) + b.get(i)); } } out.setValueCount(n); return out; } }
CREATE FUNCTION arrow_add(INT, INT) RETURNS INT PROPERTIES ( "symbol" = "com.starrocks.udf.sample.ArrowAdd", "type" = "StarrocksJar", "input" = "arrow", "file" = "http://http_host:http_port/udf-1.0-SNAPSHOT-jar-with-dependencies.jar" );

UDTF——process(FieldVector...)返回T[][](每个输入行对应一个T[]输出行;只有输入是向量化的):

public class ArrowRepeat { public Integer[][] process(IntVector v) { Integer[][] out = new Integer[v.getValueCount()][]; for (int i = 0; i < v.getValueCount(); i++) { out[i] = v.isNull(i) ? new Integer[0] : new Integer[] {v.get(i), v.get(i)}; } return out; } }

UDAF——update(State, FieldVector...)接收发往同一个 state 的整批数据;create/merge/serialize/finalize保持常规(装箱)签名。

LIMITATIONS(限制)

  • Arrow-input UDAF 仅支持全局聚合排序流式聚合(整批数据属于同一个 state)。将 Arrow UDAF 用于 hashGROUP BY的查询会以明确的错误信息失败——GROUP BYUDAF 请使用默认(装箱)输入。FE 端 CreateFunctionAnalyzer.java 的注释也印证了这一限制。
  • 窗口函数("analytic" = "true")暂不支持 Arrow 输入。

SQL 类型与方法收到的 ArrowFieldVector之间的映射关系:

SQL 类型Arrow FieldVector
BOOLEANBitVector
TINYINTTinyIntVector
SMALLINTSmallIntVector
INTIntVector
BIGINTBigIntVector
FLOATFloat4Vector
DOUBLEFloat8Vector
VARCHARVarCharVector
DECIMALDecimalVector
DATEDateDayVector
DATETIMETimeStampMicroVector
ARRAYListVector
MAPMapVector
STRUCTStructVector
创建 UDAF

执行以下命令创建上文编译的 UDAF:

CREATE [GLOBAL] AGGREGATE FUNCTION MY_SUM_INT(INT) RETURNS INT PROPERTIES ( "symbol" = "com.starrocks.udf.sample.SumInt", "type" = "StarrocksJar", "file" = "http://http_host:http_port/udf-1.0-SNAPSHOT-jar-with-dependencies.jar" );

PROPERTIES 中各参数的说明与「创建标量 UDF」相同。

创建 UDWF

执行以下命令创建上文编译的 UDWF:

CREATE [GLOBAL] AGGREGATE FUNCTION MY_WINDOW_SUM_INT(Int) RETURNS Int properties ( "analytic" = "true", "symbol" = "com.starrocks.udf.sample.WindowSumInt", "type" = "StarrocksJar", "file" = "http://http_host:http_port/udf-1.0-SNAPSHOT-jar-with-dependencies.jar" );

analytic:该 UDF 是否为窗口函数,需设置为true。其他属性的说明与「创建标量 UDF」相同。

创建 UDTF

执行以下命令创建上文编译的 UDTF:

CREATE [GLOBAL] TABLE FUNCTION MY_UDF_SPLIT(string) RETURNS string properties ( "symbol" = "com.starrocks.udf.sample.UDFSplit", "type" = "StarrocksJar", "file" = "http://http_host:http_port/udf-1.0-SNAPSHOT-jar-with-dependencies.jar" );

PROPERTIES 中各参数的说明与「创建标量 UDF」相同。

Step 7:使用 UDF

创建 UDF 后,可以根据业务需求进行测试和使用。

使用标量 UDF
SELECT MY_UDF_JSON_GET('{"key":"{\"in\":2}"}', '$.key.in');
使用 UDAF
SELECT MY_SUM_INT(col1);
使用 UDWF
SELECT MY_WINDOW_SUM_INT(intcol) OVER (PARTITION BY intcol2 ORDER BY intcol3 ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) FROM test_basic;
使用 UDTF
-- 假设存在表 t1,其列 a、b、c1 的数据如下: SELECT t1.a,t1.b,t1.c1 FROM t1; > output: 1,2.1,"hello world" 2,2.2,"hello UDTF." -- 运行 MY_UDF_SPLIT() 函数。 SELECT t1.a,t1.b, MY_UDF_SPLIT FROM t1, MY_UDF_SPLIT(t1.c1); > output: 1,2.1,"hello" 1,2.1,"world" 2,2.2,"hello" 2,2.2,"UDTF."

注意

  • 上述代码中第一个MY_UDF_SPLIT是第二个MY_UDF_SPLIT(函数)所返回列的别名。
  • 不能使用AS t2(f1)为返回的表及其列指定别名。

查看 UDF

执行以下命令查看 UDF:

SHOW [GLOBAL] FUNCTIONS;

更多信息参见 SHOW FUNCTIONS。

删除 UDF

执行以下命令删除 UDF:

DROP [GLOBAL] FUNCTION <function_name>(arg_type [, ...]);

更多信息参见 DROP FUNCTION。

SQL 数据类型与 Java 数据类型映射

注意:标量 UDF、UDAF 和 UDTF 都支持嵌套的ARRAYMAPSTRUCT参数/返回类型——包括任意嵌套,例如ARRAY<ARRAY<INT>>ARRAY<MAP<INT, STRING>>MAP<INT, ARRAY<STRING>>STRUCT<a INT, b ARRAY<STRING>>ARRAY<STRUCT<a INT, b STRING>>。叶子元素类型仍必须是下表所列的标量类型之一。由于 Java 类型擦除,对于子树中不含 STRUCT 的 ARRAY/MAP 槽位,Java 方法签名只需原始类型java.util.List/java.util.Map,StarRocks 会根据 SQL 签名驱动逐行转换。而 STRUCT 槽位必须绑定到具体的 Javarecord类,以便分析器在 JNI 边界保留正式 record 类型(见下文「STRUCT 类型绑定」)。

SQL 类型Java 类型
BOOLEANjava.lang.Boolean
TINYINTjava.lang.Byte
SMALLINTjava.lang.Short
INTjava.lang.Integer
BIGINTjava.lang.Long
FLOATjava.lang.Float
DOUBLEjava.lang.Double
STRING/VARCHARjava.lang.String
DECIMAL(p, s)(DECIMAL32 / 64 / 128 / 256)java.math.BigDecimal
DATEjava.time.LocalDate
DATETIMEjava.time.LocalDateTime
ARRAYjava.util.List
Mapjava.util.Map
STRUCTUDF 作者声明的 Javarecord

注意:对于DECIMAL参数,UDF 产生的 BigDecimal 值在写回前会使用RoundingMode.HALF_UP按列声明的(precision, scale)重新缩放。若缩放后的值超出声明的(precision, scale),行为取决于会话的overflow_mode

  • OUTPUT_NULL(默认):该行写为NULL
  • REPORT_ERROR:查询以ArithmeticException中止。

STRUCT 类型绑定

STRUCT参数与返回类型必须绑定到 UDF 作者声明的 Javarecord类(JDK 14+)。映射按位置进行:

  • record 的组件数量必须与 SQLSTRUCT字段数量一致。
  • 每个组件的类型必须按位置与对应的 SQL 字段类型匹配;不强制要求组件名与字段名一致(Java 标识符无法表达所有合法的 SQL 字段名,且 CREATE FUNCTION 按位置绑定,与会话的STRUCT_CAST_BY_NAME设置无关)。
  • 任意位置都支持嵌套STRUCT:作为 record 组件、作为ARRAY元素(List<MyRecord>)、或作为MAP的 key/value(Map<String, MyRecord>)。

同样的 record 类绑定适用于标量 UDF、UDAF 和 UDTF。对于 UDAF,record 类用于update(State, ...)参数与finalize(State)返回类型;对于 UDTF,record 类用于process(...)参数与TYPE[] process(...)的元素返回类型。

标量 UDF 示例
public record Address(String street, Integer zip) {} public record AddressOut(String full, Integer region) {} public class AddrUdf { public AddressOut evaluate(Address addr) { return new AddressOut(addr.street() + " #" + addr.zip(), addr.zip() / 1000); } }
CREATE FUNCTION addr_udf(struct<street string, zip int>) RETURNS struct<`full` string, region int> PROPERTIES ( "symbol" = "com.example.AddrUdf", "type" = "StarrocksJar", "file" = "http://localhost:8080/addr_udf.jar" );
UDAF 示例
public record Item(String name, Integer qty) {} public record TopItem(String name, Long total) {} public class TopItemAgg { public static class State { // 为简洁省略序列化实现 public int serializeLength() { return 0; } } public State create() { return new State(); } public void destroy(State state) {} public void update(State state, Item item) { /* ... */ } public void serialize(State state, java.nio.ByteBuffer buf) { /* ... */ } public void merge(State state, java.nio.ByteBuffer buf) { /* ... */ } public TopItem finalize(State state) { return new TopItem("a", 0L); } }
CREATE AGGREGATE FUNCTION top_item_agg(struct<name string, qty int>) RETURNS struct<name string, total bigint> PROPERTIES ( "symbol" = "com.example.TopItemAgg", "type" = "StarrocksJar", "file" = "http://localhost:8080/top_item_agg.jar" );
UDTF 示例
public record Pair(String key, Integer value) {} public class ExplodePairs { public Pair[] process(java.util.Map<String, Integer> m) { return m.entrySet().stream() .map(e -> new Pair(e.getKey(), e.getValue())) .toArray(Pair[]::new); } }
CREATE TABLE FUNCTION explode_pairs(map<string, int>) RETURNS struct<`key` string, `value` int> PROPERTIES ( "symbol" = "com.example.ExplodePairs", "type" = "StarrocksJar", "file" = "http://localhost:8080/explode_pairs.jar" );

参数设置

可以在集群中每台 BE 节点 JVM 的be/conf/be.conf文件中配置以下环境变量以控制内存使用:

JAVA_OPTS="-Xmx12G"

仓库中 conf/be.conf 的默认配置为:

JAVA_OPTS="--add-opens=java.base/java.util=ALL-UNNAMED --add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/sun.nio.ch=ALL-UNNAMED"

即默认已包含 Arrow 向量化输入所必需的--add-opens=java.base/java.nio=ALL-UNNAMED等 JVM 模块开放选项。调整内存时请保留这些 flag;若同时使用 Arrow 输入,务必不要移除java.base/java.nio的开放选项。从源码结构看,BE 通过 java_env.cpp 与 java_udf_context.cpp 管理嵌入 JVM 的启动与函数上下文,JAVA_OPTS直接影响该 JVM 的堆内存与模块访问权限。

FAQ

Q:创建 UDF 时可以使用静态变量吗?不同 UDF 的静态变量会相互影响吗?

A:可以,编译 UDF 时可以使用静态变量。不同 UDF 的静态变量相互隔离,即使这些 UDF 包含类名完全相同的类,彼此也不会相互影响。此外,若希望在多次 UDF 执行之间共享函数实例并支持静态变量,可在CREATE FUNCTION的 PROPERTIES 中设置"isolation" = "shared",对应 FE 端 CreateFunctionAnalyzer.java 中ISOLATION_SHARED属性的解析逻辑。

【免费下载链接】starrocksThe world's fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询