ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

StarRocks Java UDF 开发实战:从标量函数到 UDAF/UDWF/UDTF 的完整指南 StarRocks Java UDF 开发实战从标量函数到 UDAF/UDWF/UDTF 的完整指南【免费下载链接】starrocksThe worlds 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/starrocksJava UDFUser-Defined Function是 StarRocks 提供给用户的扩展机制允许使用 Java 语言编译自定义函数以满足特定业务需求。本文基于官方文档与仓库源码系统讲解标量 UDF、用户自定义聚合函数UDAF、用户自定义窗口函数UDWF与用户自定义表函数UDTF的完整开发、打包、注册与使用流程同时深入剖析底层实现原理与向量化Arrow输入等进阶特性。读完本文你将能够独立创建一个 Maven 项目编写四类 Java UDF 并在 StarRocks 集群中注册与调用。Java UDF 功能概览自 StarRocksv2.2.0起用户可以编译 Java UDF 以满足特定业务需求自v3.0起StarRocks 支持全局 UDFGlobal UDF只需在相关 SQL 语句CREATE/SHOW/DROP中加入GLOBAL关键字即可。目前 StarRocks 支持以下四类 Java UDFUDF 类型全称行为特征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 之前需要满足以下条件安装 Apache Maven用于创建和编译 Java 工程。服务器安装 JDK 17。开启 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.policyFE 的 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 version1.0 encodingUTF-8? project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion groupIdorg.example/groupId artifactIdudf/artifactId version1.0-SNAPSHOT/version properties maven.compiler.source17/maven.compiler.source maven.compiler.target17/maven.compiler.target /properties dependencies dependency groupIdcom.alibaba/groupId artifactIdfastjson/artifactId version1.2.76/version /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-dependency-plugin/artifactId version2.10/version executions execution idcopy-dependencies/id phasepackage/phase goals goalcopy-dependencies/goal /goals configuration outputDirectory${project.build.directory}/lib/outputDirectory /configuration /execution /executions /plugin plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-assembly-plugin/artifactId version3.3.0/version executions execution idmake-assembly/id phasepackage/phase goals goalsingle/goal /goals /execution /executions configuration descriptorRefs descriptorRefjar-with-dependencies/descriptorRef /descriptorRefs /configuration /plugin /plugins /build /project注意maven.compiler.source与maven.compiler.target需设置为 17与 JDK 17 前置条件保持一致。仓库示例工程 contrib/java-udf/pom.xml 同样使用java.version17。Step 3编写 UDF Java 类使用 Java 语言编写 UDF。方法中的请求参数与返回参数的数据类型必须与 Step 6 中CREATE FUNCTION语句声明的类型一致并遵循下文「SQL 数据类型与 Java 数据类型映射」中的对应关系。编写标量 UDFScalar UDF标量 UDF 对单行数据进行操作并返回单个值查询结果集中每一行对应一个值。典型的内置标量函数包括UPPER、LOWER、ROUND、ABS。业务场景示例假设 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 对多行数据进行操作并返回单个值。典型的内置聚合函数包括SUM、COUNT、MAX、MIN它们聚合每个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()创建一个状态Statevoid 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.counterval; } } 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.jar和udf-1.0-SNAPSHOT-jar-with-dependencies.jar。Step 5上传 Java 工程将udf-1.0-SNAPSHOT-jar-with-dependencies.jarfat JAR上传到一个持续运行、且集群内所有 FE 和 BE 均可访问的 HTTP 服务器上然后执行以下命令部署该文件mvn deploy也可以使用 Python 搭建一个简单的 HTTP 服务器来托管 JAR 文件。注意在 Step 6 中FE 会检查包含 UDF 代码的 JAR 文件并计算校验和BE 会下载并执行该 JAR 文件。因此 HTTP 服务器必须保持在线且所有 FE、BE 节点都能访问到该 URL。Step 6在 StarRocks 中创建 UDFStarRocks 支持在两种命名空间下创建 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否是否创建全局 UDFv3.0 起支持AGGREGATE否是否创建 UDAF 或 UDWFTABLE否是否创建 UDTF。若AGGREGATE与TABLE都未指定则创建标量函数function_name是要创建的函数名可包含数据库名例如db1.my_func。若function_name包含数据库名则在对应数据库创建 UDF否则在当前数据库创建。新函数名与其参数不能与目标数据库中已有名称相同否则创建失败函数名相同但参数不同时创建成功arg_type是函数参数类型可用, ...表示多个参数。支持的数据类型见「SQL 数据类型与 Java 数据类型映射」return_type是函数返回类型支持的数据类型同上PROPERTIES是函数属性根据创建的 UDF 类型不同而不同创建标量 UDF执行以下命令创建上文编译的标量 UDFCREATE [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_nametypeUDF 类型。设置为StarrocksJar表示该 UDF 是基于 Java 的函数file下载包含 UDF 代码的 JAR 文件的 HTTP URL格式为http://http_server_ip:http_server_port/jar_package_nameisolation可选。若要在多次 UDF 执行之间共享函数实例并支持静态变量则设置为sharedinput可选。输入格式合法值为scalar默认每行一个装箱的 Java 对象与arrow向量化每个参数一个 Apache ArrowFieldVector覆盖整批数据。参见下文「向量化Arrow输入」源码佐证FE 端 CreateFunctionAnalyzer.java 负责解析CREATE FUNCTION语句读取symbol、input、isolation等属性如CreateFunctionStmt.SYMBOL_KEY、ISOLATION_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 设为providedBE 已自带该依赖dependency groupIdorg.apache.arrow/groupId artifactIdarrow-vector/artifactId version17.0.0/version scopeprovided/scope /dependency结果向量请通过arg.getAllocator()从框架托管的 allocator 中分配这样引擎可以接管其生命周期。注意BE 在嵌入式 JVM 中运行 UDF该 JVM 必须向 Arrow 的堆外内存层暴露java.nio。仓库自带的 conf/be.conf 已通过JAVA_OPTS... --add-opensjava.base/java.nioALL-UNNAMED ...完成该设置如果你自定义JAVA_OPTS必须保留该 flag否则 arrow-input UDF 将初始化失败。标量 UDF——evaluate(FieldVector...)返回一个FieldVectorpublic 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 FieldVectorBOOLEANBitVectorTINYINTTinyIntVectorSMALLINTSmallIntVectorINTIntVectorBIGINTBigIntVectorFLOATFloat4VectorDOUBLEFloat8VectorVARCHARVarCharVectorDECIMALDecimalVectorDATEDateDayVectorDATETIMETimeStampMicroVectorARRAYListVectorMAPMapVectorSTRUCTStructVector创建 UDAF执行以下命令创建上文编译的 UDAFCREATE [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执行以下命令创建上文编译的 UDWFCREATE [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执行以下命令创建上文编译的 UDTFCREATE [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 后可以根据业务需求进行测试和使用。使用标量 UDFSELECT MY_UDF_JSON_GET({key:{\in\:2}}, $.key.in);使用 UDAFSELECT MY_SUM_INT(col1);使用 UDWFSELECT 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执行以下命令查看 UDFSHOW [GLOBAL] FUNCTIONS;更多信息参见 SHOW FUNCTIONS。删除 UDF执行以下命令删除 UDFDROP [GLOBAL] FUNCTION function_name(arg_type [, ...]);更多信息参见 DROP FUNCTION。SQL 数据类型与 Java 数据类型映射注意标量 UDF、UDAF 和 UDTF 都支持嵌套的ARRAY、MAP和STRUCT参数/返回类型——包括任意嵌套例如ARRAYARRAYINT、ARRAYMAPINT, STRING、MAPINT, ARRAYSTRING、STRUCTa INT, b ARRAYSTRING、ARRAYSTRUCTa INT, b STRING。叶子元素类型仍必须是下表所列的标量类型之一。由于 Java 类型擦除对于子树中不含 STRUCT 的 ARRAY/MAP 槽位Java 方法签名只需原始类型java.util.List/java.util.MapStarRocks 会根据 SQL 签名驱动逐行转换。而 STRUCT 槽位必须绑定到具体的 Javarecord类以便分析器在 JNI 边界保留正式 record 类型见下文「STRUCT 类型绑定」。SQL 类型Java 类型BOOLEANjava.lang.BooleanTINYINTjava.lang.ByteSMALLINTjava.lang.ShortINTjava.lang.IntegerBIGINTjava.lang.LongFLOATjava.lang.FloatDOUBLEjava.lang.DoubleSTRING/VARCHARjava.lang.StringDECIMAL(p, s)DECIMAL32 / 64 / 128 / 256java.math.BigDecimalDATEjava.time.LocalDateDATETIMEjava.time.LocalDateTimeARRAYjava.util.ListMapjava.util.MapSTRUCTUDF 作者声明的 Javarecord类注意对于DECIMAL参数UDF 产生的 BigDecimal 值在写回前会使用RoundingMode.HALF_UP按列声明的(precision, scale)重新缩放。若缩放后的值超出声明的(precision, scale)行为取决于会话的overflow_modeOUTPUT_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元素ListMyRecord、或作为MAP的 key/valueMapString, MyRecord。同样的 record 类绑定适用于标量 UDF、UDAF 和 UDTF。对于 UDAFrecord 类用于update(State, ...)参数与finalize(State)返回类型对于 UDTFrecord 类用于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(structstreet string, zip int) RETURNS structfull 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(structname string, qty int) RETURNS structname 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.MapString, Integer m) { return m.entrySet().stream() .map(e - new Pair(e.getKey(), e.getValue())) .toArray(Pair[]::new); } }CREATE TABLE FUNCTION explode_pairs(mapstring, int) RETURNS structkey 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-opensjava.base/java.utilALL-UNNAMED --add-opensjava.base/java.nioALL-UNNAMED --add-opensjava.base/sun.nio.chALL-UNNAMED即默认已包含 Arrow 向量化输入所必需的--add-opensjava.base/java.nioALL-UNNAMED等 JVM 模块开放选项。调整内存时请保留这些 flag若同时使用 Arrow 输入务必不要移除java.base/java.nio的开放选项。从源码结构看BE 通过 java_env.cpp 与 java_udf_context.cpp 管理嵌入 JVM 的启动与函数上下文JAVA_OPTS直接影响该 JVM 的堆内存与模块访问权限。FAQQ创建 UDF 时可以使用静态变量吗不同 UDF 的静态变量会相互影响吗A可以编译 UDF 时可以使用静态变量。不同 UDF 的静态变量相互隔离即使这些 UDF 包含类名完全相同的类彼此也不会相互影响。此外若希望在多次 UDF 执行之间共享函数实例并支持静态变量可在CREATE FUNCTION的 PROPERTIES 中设置isolation shared对应 FE 端 CreateFunctionAnalyzer.java 中ISOLATION_SHARED属性的解析逻辑。【免费下载链接】starrocksThe worlds 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),仅供参考
返回列表