- OLAP
- 数据库
- 大数据
- 实时分析
【免费下载链接】doris
Apache Doris is an easy-to-use, high performance and unified analytics database.
Apache Doris 提供了 RPC(Remote Procedure Call)形式的远程 UDF 能力:函数逻辑运行在独立的 RPC 服务进程中,通过 gRPC/brpc 协议与 Doris 集群交互。本文以仓库中的 remote-udf-java-demo 示例为主线,完整讲解如何用 Java 构建一个 Remote UDF Function Service——从 proto 协议定义、服务端实现到编译运行,并结合 BE 端源码剖析 RPC 调用的底层链路。读完本文,你将能够独立编写并部署自己的 Java 版远程函数服务。
一、Remote UDF 是什么:函数逻辑外置的架构思路
传统的 Doris UDF 以动态库(.so)形式加载进 BE 进程内执行。Remote UDF 则采用相反的设计:将函数实现部署为一个独立的 RPC 服务(Function Service),BE 在查询执行时通过网络把批量参数打包发送给该服务,由服务端计算后返回结果列。这样做的好处是:
- 语言无关:只要实现约定的 proto 接口,可以用任意语言(Java、C++、Python 等)编写函数;
- 故障隔离:函数崩溃不会拖垮 BE 进程,服务可独立扩容、独立发布;
- 资源解耦:计算密集或依赖外部系统的函数可放在专用服务上运行。
在 Doris 中,这类远程函数被称作 RPC UDF,其核心接口定义在仓库根目录的 gensrc/proto/function_service.proto,与示例中使用的 proto 文件内容一致。
二、Demo 项目结构与构建方式
2.1 目录与文件总览
Demo 项目位于 samples/doris-demo/remote-udf-java-demo,是一个标准 Maven Java 工程:
samples/doris-demo/remote-udf-java-demo/ ├── pom.xml # Maven 构建配置 ├── README.md # 编译与运行说明 └── src/main/ ├── java/org/apache/doris/udf/ │ ├── FunctionServiceDemo.java # gRPC 服务启动入口 │ └── FunctionServiceImpl.java # 三个 RPC 方法的实现 └── proto/ ├── function_service.proto # RPC 服务与消息定义 └── types.proto # 数据类型与通用消息定义2.2 Maven 依赖与构建配置(pom.xml)
pom.xml 中声明的关键依赖如下:
<properties> <protobuf.version>3.15.0</protobuf.version> <grpc.version>1.44.1</grpc.version> <java.version>1.8</java.version> </properties> <dependencies> <dependency> <groupId>io.grpc</groupId> <artifactId>grpc-protobuf</artifactId> <version>${grpc.version}</version> </dependency> <dependency> <groupId>io.grpc</groupId> <artifactId>grpc-stub</artifactId> <version>${grpc.version}</version> </dependency> <dependency> <groupId>io.grpc</groupId> <artifactId>grpc-netty-shaded</artifactId> <version>${grpc.version}</version> </dependency> </dependencies>构建链路中有三个值得注意的细节:
- proto 代码生成:使用
protoc-jar-maven-plugin在generate-sources阶段从src/main/proto编译生成 Java 与 gRPC-Java 桩代码。插件默认使用${basedir}/../../../thirdparty/installed/bin/protoc(即仓库 thirdparty 目录下安装的 protoc)作为编译器;pom 中注释也给出了备选方案:改用protocArtifact(如com.google.protobuf:protoc:3.15.0)后即可免去本地安装 protobuf 工具链。 - 生成源码注册:
build-helper-maven-plugin把target/generated-sources加入编译源码目录,保证生成的org.apache.doris.proto包类可被引用。 - 可执行 Jar:
maven-assembly-plugin的jar-with-dependencies描述符将全部依赖打进单个 fat jar,maven-jar-plugin与 assembly 均指定主类为org.apache.doris.udf.FunctionServiceDemo。
三、协议定义:一次远程函数调用长什么样
协议层由两个 proto 文件组成,它们是 BE 端 C++ 客户端与 Java 服务端之间的"契约"。示例内的 function_service.proto 与仓库根目录 gensrc/proto/function_service.proto 保持同构(proto2语法,java_package = "org.apache.doris.proto")。
3.1 PFunctionService:三个 RPC 方法
service PFunctionService { rpc fn_call(PFunctionCallRequest) returns (PFunctionCallResponse); rpc check_fn(PCheckFunctionRequest) returns (PCheckFunctionResponse); rpc hand_shake(PHandShakeRequest) returns (PHandShakeResponse); }fn_call:核心调用入口。BE 把一列参数打包进PFunctionCallRequest,服务端执行函数逻辑并返回结果列PFunctionCallResponse;check_fn:函数校验。BE 在注册/执行前校验函数签名(函数名、参数个数等)是否匹配;hand_shake:握手。用于连接建立阶段的连通性与协议版本探测,请求携带一个hello字符串,响应原样带回。
3.2 请求与响应消息
message PFunctionCallRequest { optional string function_name = 1; repeated PValues args = 2; optional PRequestContext context = 3; } message PFunctionCallResponse { repeated PValues result = 1; optional PStatus status = 2; optional PRequestContext context = 3; }function_name即 UDF 注册时使用的函数名(在 BE 侧它对应TFunction.scalar_fn.symbol);args是按参数位置排列的批量值列(PValues),天然支持向量化的一次性多行计算。
3.3 数据类型:PValues 与 PGenericType
types.proto 定义了跨 RPC 传输的类型系统,核心是PValues(一列批量值)与PGenericType(类型描述):
message PValues { required PGenericType type = 1; optional bool has_null = 2 [default = false]; repeated bool null_map = 3; repeated double double_value = 4; repeated float float_value = 5; repeated int32 int32_value = 6; repeated int64 int64_value = 7; repeated uint32 uint32_value = 8; repeated uint64 uint64_value = 9; repeated bool bool_value = 10; repeated string string_value = 11; repeated bytes bytes_value = 12; repeated PDateTime datetime_value = 13; repeated PValues child_element = 14; // 复杂类型(如 ARRAY)的子元素 repeated int64 child_offset = 15; }PGenericType.TypeId枚举覆盖了 Doris 的主要类型,例如INT32、INT64、DOUBLE、STRING、DATEV2、DATETIMEV2、DECIMAL128、JSONB、VARIANT、IPV4/IPV6等,并对LIST、MAP、STRUCT复杂类型通过PList/PMap/PStruct/PDecimal描述子类型。PStatus通过status_code(0 表示成功)和error_msgs汇报错误。
四、服务端实现:三个示例函数逐一拆解
FunctionServiceImpl.java 继承 gRPC 生成的PFunctionServiceGrpc.PFunctionServiceImplBase,实现了三个方法,并内置了三个可演示的函数。
4.1 fn_call:按函数名分发执行
@Override public void fnCall(FunctionService.PFunctionCallRequest request, StreamObserver<FunctionService.PFunctionCallResponse> responseObserver) { String functionName = request.getFunctionName(); FunctionService.PFunctionCallResponse res = null; if ("add_int_two".equals(functionName)) { res = FunctionService.PFunctionCallResponse.newBuilder() .setStatus(Types.PStatus.newBuilder().setStatusCode(0).build()) .addResult(Types.PValues.newBuilder().setHasNull(false) .addAllInt32Value(IntStream.range(0, Math.min(request.getArgs(0) .getInt32ValueCount(), request.getArgs(1).getInt32ValueCount())) .mapToObj(i -> request.getArgs(0).getInt32Value(i) + request.getArgs(1).getInt32Value(i)).collect(Collectors.toList())) .setType(Types.PGenericType.newBuilder() .setId(Types.PGenericType.TypeId.INT32).build()) .build()).build(); } // add_int_one / add_string 分支省略,逻辑同理 ok(responseObserver, res); }示例提供了三个函数:
| 函数名 | 输入 | 输出 | 计算逻辑 |
|---|---|---|---|
add_int_two | 两列 INT32 | INT32 | 逐行相加,行数取两列最小值 |
add_int_one | 一列 INT32 | INT32 | 每个元素加 1 |
add_string | 一列 STRING | STRING | 每个元素拼接_rpc_test后缀 |
注意实现完全面向批量:从args(i)中取出整列int32_value/string_value列表,逐元素计算后以addAllInt32Value/addAllStringValue一次性写回,体现向量化执行风格。
4.2 check_fn:调用前的参数校验
@Override public void checkFn(FunctionService.PCheckFunctionRequest request, StreamObserver<FunctionService.PCheckFunctionResponse> responseObserver) { int status = 0; if ("add_int_two".equals(request.getFunction().getFunctionName())) { if (request.getFunction().getInputsCount() != 2) { // 校验参数个数 status = -1; } } // add_int_one / add_string 分支校验 inputsCount 是否为 1 FunctionService.PCheckFunctionResponse res = FunctionService.PCheckFunctionResponse.newBuilder() .setStatus(Types.PStatus.newBuilder().setStatusCode(status).build()).build(); ok(responseObserver, res); }校验逻辑很简单:根据函数名核对输入参数个数,不匹配时status_code置为 -1。实际项目中可在此基础上扩展类型匹配、常量折叠等更严格的检查。
4.3 handShake:连通性握手
@Override public void handShake(Types.PHandShakeRequest request, StreamObserver<Types.PHandShakeResponse> responseObserver) { ok(responseObserver, Types.PHandShakeResponse.newBuilder() .setStatus(Types.PStatus.newBuilder().setStatusCode(0).build()) .setHello(request.getHello()).build()); }把请求中的hello原样返回,同时返回成功状态码,用于验证服务存活与协议兼容。
五、服务启动入口与端口参数
FunctionServiceDemo.java 负责拉起 gRPC 服务:
private void start(int port) throws IOException { server = ServerBuilder.forPort(port) .addService(new FunctionServiceImpl()) .build() .start(); logger.info("Server started, listening on " + port); // 注册 JVM 关闭钩子,优雅停机 Runtime.getRuntime().addShutdownHook(new Thread(() -> { System.err.println("*** shutting down gRPC server since JVM is shutting down"); try { FunctionServiceDemo.this.stop(); } catch (InterruptedException e) { e.printStackTrace(System.err); } System.err.println("*** server shut down"); })); } public static void main(String[] args) throws IOException, InterruptedException { int port = 9000; if (args.length > 0) { port = Integer.parseInt(args[0]); // 端口必须为正整数 } if (port <= 0) { System.err.println("port " + args[0] + " must be positive."); System.exit(1); } final FunctionServiceDemo server = new FunctionServiceDemo(); server.start(port); server.blockUntilShutdown(); // 主线程阻塞等待 }实现要点:端口号通过命令行第一个参数传入(默认 9000),非整数或非正数会直接退出;stop()使用shutdown().awaitTermination(30, TimeUnit.SECONDS)等待在途请求最多 30 秒后优雅退出。
六、编译与运行(严格按原文档步骤执行)
README 给出了两步核心操作,完整展开如下。
6.1 编译打包
mvn package执行时generate-sources阶段会用仓库 thirdparty 下的 protoc 编译两个 proto 文件并生成 gRPC 桩代码,随后编译 Java 源码并打两个 jar:
target/remote-udf-java-demo.jar:仅含项目自身类;target/remote-udf-java-demo-jar-with-dependencies.jar:内含全部依赖的 fat jar(默认运行方式)。
若本地没有 Doris 的 thirdparty 工具链,可按 pom 中注释修改为protocArtifact方式,让 Maven 自动下载 protoc。
6.2 启动服务
java -jar target/remote-udf-java-demo-jar-with-dependencies.jar 9000其中9000是服务监听端口(可按需更换为任意正整数)。启动成功后可看到日志Server started, listening on 9000。由于 gRPC 使用 daemon 线程,主线程通过blockUntilShutdown()阻塞,保证进程常驻等待查询请求。
6.3 在 Doris 中注册并使用
服务运行后,在 Doris 中通过 SQL 将该服务声明为远程 UDF(示意,参数以当前集群的 FE 语法为准):
CREATE FUNCTION add_int_two(int, int) RETURNS int PROPERTIES ( "type" = "RPC", "symbol" = "add_int_two", "object_file" = "127.0.0.1:9000" ); SELECT add_int_two(1, 2);其中symbol对应服务端request.getFunctionName()匹配的函数名,object_file指向服务地址。
七、BE 端底层链路:从 SQL 到 RPC 的源码印证
为验证协议与调用约定,下面从 BE 源码看一次远程调用的完整路径。
7.1 调用链总览
- SQL 中的 RPC 函数经 FE 解析后,在 BE 中被实例化为
FunctionRPC(be/src/vec/functions/function_rpc.h),其open()在FRAGMENT_LOCAL作用域创建RPCFnImpl并持有 brpc stub 客户端; execute()委托RPCFnImpl::vec_call()执行:- 用
_convert_block_to_proto把查询执行中的Block参数列序列化为PFunctionCallRequest; - 调用
_client->fn_call(...)发起 brpc 调用(见 be/src/vec/functions/function_rpc.cpp); - 检查响应
status_code == 0后,用_convert_to_block把返回的PValues反序列化回Block的对应列。
- 用
7.2 关键实现细节
Status RPCFnImpl::vec_call(FunctionContext* context, Block& block, const ColumnNumbers& arguments, size_t result, size_t input_rows_count) { PFunctionCallRequest request; PFunctionCallResponse response; request.set_function_name(_function_name); // 对应 Java 端的 functionName RETURN_IF_ERROR(_convert_block_to_proto(block, arguments, input_rows_count, &request)); brpc::Controller cntl; _client->fn_call(&cntl, &request, &response, nullptr); // 阻塞式 RPC if (cntl.Failed()) { /* ... */ } if (!response.has_status() || response.result_size() == 0) { /* ... */ } if (response.status().status_code() != 0) { /* ... */ } _convert_to_block(block, response.result(0), result); // 取第一个结果列 return Status::OK(); }几点值得注意:
- 函数名映射:BE 侧
_function_name = _fn.scalar_fn.symbol,即注册 UDF 时symbol属性,与 Java 端request.getFunctionName()匹配,因此示例中的函数名add_int_two等必须与注册时的 symbol 完全一致; - 服务地址映射:
_server_addr = _fn.hdfs_location取自注册时的object_file(即host:port),客户端通过ExecEnv的 brpc 客户端缓存按地址复用连接(be/src/util/brpc_client_cache.cpp); - 参数打包:
_convert_block_to_proto逐列把Block序列化为PValues追加进request.args,has_null标记整列是否含 NULL,常量列通过convert_to_full_column_if_const展开为全量列; - 结果回填:
_convert_to_block只取response.result(0),把PValues通过DataTypeSerde::read_column_from_pb反序列化后替换目标位置列。
也就是说,Java 服务端返回的PValues类型必须与注册函数的返回类型严格对应(例如add_int_two返回TypeId.INT32),否则反序列化阶段会失败。
八、编写自定义远程函数的注意事项
结合协议定义与两端实现,总结实践要点:
- 函数名即路由键:
fn_call通过function_name分发,服务端需自行实现 if/else 或注册表式路由,并与 Doris 注册的symbol保持一致; - 返回值类型必须精确匹配:返回的
PValues.type(如INT32、STRING)需与函数声明的返回类型一致,NULL 支持通过has_null与null_map表达; - check_fn 尽量严格:参数个数/类型校验失败应返回非 0 状态码,便于在查询规划期就暴露问题,而不是执行期才报错;
- 保持向量化:请求按整列批量下发,服务端应整列读取、整列返回,避免逐行 RPC;
- 版本兼容:proto 协议若与当前 BE 版本不一致(新增字段、枚举值),建议以仓库根目录 gensrc/proto/function_service.proto 与 gensrc/proto/types.proto 为唯一事实来源同步更新;
- 部署形态:服务是独立进程,需保证 BE 到服务地址的网络连通性,并配合进程守护工具常驻运行。
九、小结
本文以仓库自带的 Java Demo 为骨架,完整走通了 Remote UDF Function Service 的协议定义、服务实现、编译运行与 BE 调用链路:PFunctionService的三个 RPC 方法(fn_call/check_fn/hand_shake)构成了远程函数的通信契约,PValues批量值结构保证了向量化传输,而 BE 端FunctionRPC(be/src/vec/functions/function_rpc.cpp)则负责列数据与 proto 消息的双向转换。掌握这套模式后,你可以基于 samples/doris-demo/remote-udf-java-demo 快速扩展出任意自定义远程函数,甚至用其他语言按同一份 proto 实现同样的服务。
- OLAP
- 数据库
- 大数据
- 实时分析
【免费下载链接】doris
Apache Doris is an easy-to-use, high performance and unified analytics database.
相关推荐
OpenProject 多语言配置实操指南:让 50 种语言团队各用母语
OpenProject 多语言配置实操指南:让 50 种语言团队各用母语 OpenProject 是一款开源项目管理软件,支持 50 多种界面语言,每位成员都能
OLAP数据库大数据实时分析Apache Doris Java UDF开发指南:5步实现业务逻辑高效扩展
Apache Doris Java UDF开发指南:5步实现业务逻辑高效扩展 Apache Doris作为一款高性能的统一分析数据库,其Java UDF(用户定
OLAP数据库大数据实时分析Apache Doris自定义函数开发终极指南:C++ UDF实战教程
Apache Doris作为一款高性能的统一分析数据库,其强大的自定义函数(UDF)功能让开发者能够扩展数据库的核心能力。本文将为您详细介绍如何使用C++开发A
OLAP数据库大数据实时分析
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考