Skip to content

内置 UDF

本文档介绍 SQLRec 内置的用户定义函数(UDF),包括表函数(Table Function)和标量函数(Scalar Function)。

概述

SQLRec 提供了丰富的内置 UDF,用于推荐系统开发中的常见操作,如去重、打散、向量计算等。

UDF 分类

类型说明返回值
Table Function表函数,接收表作为输入,返回表CacheTable
Scalar Function标量函数,接收标量值,返回标量值单个值

Calcite 内置函数

除 SQLRec 自身提供的 UDF 外,普通 SELECT 查询还可以直接使用 Apache Calcite 的标准内置函数。这些函数来自 Calcite 1.32.0 的 SqlStdOperatorTable,无需通过 CREATE FUNCTION 注册。

sql
SELECT ABS(-10), POWER(2, 3), UPPER('sqlrec');

SELECT category,
       COUNT(*) AS item_count,
       ROW_NUMBER() OVER (PARTITION BY category ORDER BY score DESC) AS row_num
FROM candidates
GROUP BY category, score;

SELECT JSON_VALUE('{"item_id": 1001}', 'strict $.item_id');

SQLRec 当前已经验证支持的主要函数如下。表格中的函数名是代表性清单,同一函数支持的参数类型和具体语法请参考 Calcite SQL 文档。

分类函数
数值函数ABSACOSASINATANATAN2CBRTCEILCOSCOTDEGREESEXPFLOORLNLOG10MODPIPOWERRADIANSRANDRAND_INTEGERROUNDSIGNSINSQRTTANTRUNCATE
字符串函数ASCIICHAR_LENGTHCHARACTER_LENGTHUPPERLOWERINITCAPPOSITIONOVERLAYREPLACESUBSTRINGTRIM,以及字符串连接运算符 ||
空值和类型转换COALESCENULLIFCASTCASE
日期时间函数CURRENT_DATECURRENT_TIMECURRENT_TIMESTAMPEXTRACTYEARQUARTERMONTHWEEKDAYOFYEARDAYOFMONTHDAYOFWEEKHOURMINUTESECONDLAST_DAYTIMESTAMPADDTIMESTAMPDIFF
集合函数ARRAYMAPMULTISET 构造器,集合下标访问,以及 CARDINALITYELEMENT
聚合函数COUNTSUMAVGMINMAXEVERYSOMEANY_VALUESINGLE_VALUEMODEAPPROX_COUNT_DISTINCTBIT_ANDBIT_ORBIT_XORCOLLECTLISTAGGFUSIONINTERSECTION
统计聚合函数STDDEVSTDDEV_POPSTDDEV_SAMPVARIANCEVAR_POPVAR_SAMPCOVAR_POPCOVAR_SAMPREGR_COUNTREGR_SXXREGR_SYY
分组和窗口函数GROUPINGGROUPING_IDGROUP_IDROW_NUMBERRANKDENSE_RANKNTILEFIRST_VALUELAST_VALUENTH_VALUELEADLAG
JSON 函数JSON_EXISTSJSON_VALUEJSON_QUERYJSON_OBJECTJSON_ARRAYJSON_OBJECTAGGJSON_ARRAYAGGJSON_TYPEJSON_DEPTHJSON_LENGTHJSON_KEYSJSON_PRETTYJSON_REMOVEJSON_STORAGE_SIZE

使用时请注意以下限制:

  • PI 是无参数函数,语法为 SELECT PI,而不是 PI()
  • SQLRec 当前只接入 Calcite 标准函数表,没有接入 MySQL、PostgreSQL、Oracle、Spark 或 BigQuery 等 SqlLibrary 方言扩展。使用 MySQL 风格词法解析并不表示 IFNULLDATE_FORMATNVL 等方言函数自动可用。
  • CUME_DISTPERCENT_RANKPERCENTILE_CONTPERCENTILE_DISC 当前无法生成可执行的 Enumerable 计划。
  • LOCALTIMELOCALTIMESTAMPUSERCURRENT_USERSESSION_USERSYSTEM_USERCURRENT_SCHEMA 依赖的运行上下文尚未完整提供。
  • OCTET_LENGTH 对普通 VARCHAR 不可用;Calcite 的当前实现要求二进制字符串。
  • TUMBLEHOPSESSION 以及 MATCH_RECOGNIZE 专用函数不属于普通标量 UDF,目前没有 SQLRec 端到端支持保证。

相关参考:

Calcite 官方参考页可能对应比 SQLRec 依赖更新的版本,并且包含方言扩展。判断函数是否可用于 SQLRec 时,应以上述已验证清单和项目实际依赖版本为准。

表函数(Table Function)

dedup

去重函数,根据指定列从输入表中排除已存在于去重表中的记录。

函数签名

java
public CacheTable evaluate(CacheTable input, CacheTable dedupTable, String col1, String col2)

参数说明

参数类型说明
inputCacheTable输入表
dedupTableCacheTable去重表,包含需要排除的值
col1String输入表中用于去重的列名
col2String去重表中用于匹配的列名

返回值:返回去重后的 CacheTable,结构与输入表相同。

使用示例

sql
-- 获取用户已曝光的物品
CACHE TABLE exposured_item AS
SELECT item_id
FROM user_info JOIN exposure_item ON user_id = user_info.id;

-- 从召回结果中排除已曝光物品
CACHE TABLE dedup_recall AS
CALL dedup(recall_item, exposured_item, 'item_id', 'item_id');

工作原理

  1. dedupTablecol2 列收集所有值
  2. 遍历 input 表,排除 col1 列值存在于去重集合中的记录
  3. 返回去重后的结果表

shuffle

随机打乱函数,将输入表中的记录随机排序。

函数签名

java
public CacheTable evaluate(CacheTable input)

参数说明

参数类型说明
inputCacheTable输入表

返回值:返回随机排序后的 CacheTable,结构和数据与输入表相同。

使用示例

sql
-- 随机打乱推荐结果
CACHE TABLE shuffled_result AS
CALL shuffle(recall_item);

-- 取打乱后的前 N 个
CACHE TABLE random_top_n AS
SELECT * FROM shuffled_result LIMIT 10;

window_diversify

窗口打散函数,确保相邻的记录不会过于集中在某个类目,实现推荐结果的多样性。

函数签名

java
public CacheTable evaluate(
    CacheTable input,
    String categoryColumnName,
    String windowSize,
    String maxCategoryNumInWindow,
    String maxReturnRecord
)

参数说明

参数类型说明
inputCacheTable输入表
categoryColumnNameString类目列名,用于打散的依据
windowSizeString滑动窗口大小
maxCategoryNumInWindowString窗口内每个类目最多出现的次数
maxReturnRecordString最大返回记录数

返回值:返回打散后的 CacheTable,结构与输入表相同。

使用示例

sql
-- 类目打散:窗口大小为 3,每个类目在窗口内最多出现 1 次,返回 10 条
CACHE TABLE diversify_result AS
CALL window_diversify(rec_item, 'category1', '3', '1', '10');

工作原理

  1. 维护一个滑动窗口,统计窗口内各类目的出现次数
  2. 遍历输入记录,优先选择窗口内未超限的类目
  3. 当窗口滑动时,移除最旧记录的类目计数
  4. 确保推荐结果的多样性,避免同类目物品连续出现

add_col

添加列函数,为输入表添加一个新列,所有行的该列值相同。

函数签名

java
public CacheTable evaluate(CacheTable input, String colName, String value)

参数说明

参数类型说明
inputCacheTable输入表
colNameString新列名
valueString新列的值(所有行相同)

返回值:返回添加新列后的 CacheTable

使用示例

sql
-- 添加一个来源标识列
CACHE TABLE result_with_source AS
CALL add_col(recall_item, 'source', 'daily_rec');

-- 添加时间戳列
CACHE TABLE result_with_time AS
CALL add_col(recall_item, 'rec_time', '2024-01-01');

注意事项

  • 新列名不能与已有列名重复
  • 新列类型为 VARCHAR

call_service

模型服务调用函数,用于调用已部署的模型服务进行推理。详见 模型文档

函数签名

java
public CacheTable evaluate(ReadonlyContext context, String serviceName, CacheTable input)

public CacheTable evaluate(ReadonlyContext context, String serviceName,
                           CacheTable user, CacheTable item)

context 由 SQLRec 自动注入。SQL 调用支持以下两种形式:

  • CALL call_service(serviceName, input):使用行式 JSON 调用服务,返回表包含输入列以及模型输出列。
  • CALL call_service(serviceName, user, item):使用 User-Item 模式调用服务,请求体为列式 JSON,返回表保留 Item 表列并追加模型输出列。

User-Item 模式的请求规则:

  1. User 表必须恰好一行;Item 表可以包含多行。
  2. 仅处理模型定义的输入字段。字段名如果存在于 User 表中(匹配时忽略大小写),该字段归入 User;否则归入 Item。若两张表存在同名字段,优先使用 User 表中的字段;在对应表中找不到的模型字段不会写入请求。两张表中不属于模型输入的额外列也不会写入请求。
  3. 每个字段都序列化为一个 JSON 数组。User 字段是单元素数组,因此一份用户数据在一次请求中只发送一遍;Item 字段按 Item 表的行顺序组成数组,不会为每个 Item 重复发送 User 数据。

例如,一行 User 数据和三行 Item 数据会生成:

json
{
  "user_id": [1001],
  "user_age": [25],
  "item_id": [1, 2, 3],
  "category": ["phone", "tablet", "laptop"]
}

当 Item 表为空时不会发送 HTTP 请求,函数直接返回空表,其字段为 Item 表字段加模型输出字段。

HTTP 请求超时默认为连接、读写各 30 秒。详细协议和示例参见模型文档


batch_call_service

批量模型服务调用函数,用于在 Flink SQL 中批量调用已部署的模型服务进行推理。该函数将多行数据批量发送到远程服务,并将返回结果与原始数据合并输出。

注意

此函数只能在 Flink SQL 中使用,不支持 SQLRec 的 CACHE TABLE 语法。

该函数是 Flink TableFunction,通过 LATERAL TABLE 调用;每条输入记录先进入缓冲区,达到 batchSize 后批量 POST JSON 数组,算子关闭时还会刷新不足一批的记录。服务响应必须是 JSON 对象;数组值按输入行序映射。

函数签名

java
@FunctionHint(output = @DataTypeHint("ROW<" +
        "long_map MAP<STRING, BIGINT>, " +
        "double_map MAP<STRING, DOUBLE>, " +
        "string_map MAP<STRING, STRING>, " +
        "long_array_map MAP<STRING, ARRAY<BIGINT>>, " +
        "double_array_map MAP<STRING, ARRAY<DOUBLE>>, " +
        "string_array_map MAP<STRING, ARRAY<STRING>>" +
        ">"))
public void eval(@DataTypeHint(inputGroup = InputGroup.ANY) Object... args)

参数说明

参数类型说明
serviceUrlString模型服务的 URL 地址
batchSizeInteger批量大小,每次请求发送的行数
fieldName-value pairsObject...字段名-值对,用于指定要发送到服务的字段,必须成对出现

返回值:返回一个 ROW 类型,包含以下字段:

字段名类型说明
long_mapMAP<STRING, BIGINT>长整型字段的 Map
double_mapMAP<STRING, DOUBLE>双精度浮点型字段的 Map
string_mapMAP<STRING, STRING>字符串型字段的 Map
long_array_mapMAP<STRING, ARRAY<BIGINT>>长整型数组字段的 Map
double_array_mapMAP<STRING, ARRAY<DOUBLE>>双精度浮点型数组字段的 Map
string_array_mapMAP<STRING, ARRAY<STRING>>字符串型数组字段的 Map

使用示例

sql
-- 创建临时函数
CREATE TEMPORARY FUNCTION batch_call_service AS 'com.sqlrec.udf.udtf.BatchCallServiceUDTF';

-- 调用模型服务生成物品向量
INSERT INTO item_embedding
SELECT 
    r.long_map['movie_id'] AS id,
    r.string_map['title'] AS title,
    r.string_array_map['genres'] AS genres,
    r.double_array_map['item_tower_emb'] AS embedding
FROM ml_movies, LATERAL TABLE(batch_call_service(
    'http://test-recall-service-item.sqlrec.svc.cluster.local:80/predict', 
    128, 
    'movie_id', movie_id, 
    'title', title, 
    'genres', genres
)) AS r
WHERE dt = '2024-01-01';

工作原理

  1. 函数接收多行数据,将字段名-值对缓存到缓冲区
  2. 当缓冲区大小达到 batchSize 时,将数据批量发送到模型服务
  3. 模型服务接收 JSON 数组格式的请求,返回包含预测结果的 JSON 对象
  4. 函数将预测结果与原始数据合并,按类型分类存储到不同的 Map 中
  5. 每行数据输出一个 ROW,可通过 Map 访问原始字段和预测结果

请求格式

发送到模型服务的 JSON 格式为对象数组:

json
[
  {"movie_id": 1, "title": "Toy Story", "genres": ["Animation", "Comedy"]},
  {"movie_id": 2, "title": "Jumanji", "genres": ["Adventure", "Children"]}
]

响应格式

模型服务应返回一个 JSON 对象,其中每个字段的值是一个数组,数组长度与请求数据行数相同:

json
{
  "item_tower_emb": [[0.1, 0.2, ...], [0.3, 0.4, ...]],
  "score": [0.95, 0.87]
}

注意事项

  • 此函数只能在 Flink SQL 中使用,需要使用 LATERAL TABLE 语法
  • batchSize 建议根据模型服务的性能和网络延迟进行调整,通常设置为 64-256
  • 模型服务需要支持 POST 请求,接收 JSON 数组并返回 JSON 对象
  • 返回结果中的数组字段会自动按行索引与输入数据对应
  • close() 方法中会处理缓冲区中剩余的数据

dpp_diversity

DPP(Determinantal Point Process)多样性函数,基于行列式点过程实现推荐结果的多样性打散。采用快速贪心 MAP 推理算法(参考 Hulu NIPS 2018 论文),在保证相关性的同时提升推荐结果的多样性。

函数签名

java
public CacheTable evaluate(
    CacheTable input,
    String embeddingColumnName,
    String scoreColumnName,
    String theta,
    String maxLength
)

参数说明

参数类型说明
inputCacheTable输入表
embeddingColumnNameString向量列名,用于计算物品间的相似度
scoreColumnNameString相关性得分列名,用于衡量物品质量
thetaString相关性-多样性权衡参数,取值范围 [0, 1),值越接近 1 越偏向相关性,越接近 0 越偏向多样性
maxLengthString最大返回记录数

返回值:返回多样性选择后的 CacheTable,结构与输入表相同。

使用示例

sql
-- DPP 多样性打散:theta=0.5 平衡相关性与多样性,返回 20 条
CACHE TABLE dpp_result AS
CALL dpp_diversity(rec_item, 'item_embedding', 'score', '0.5', '20');

工作原理

  1. 从输入表中提取相关性得分和向量
  2. 对负得分裁剪为极小正值,然后进行指数变换:score = exp(alpha * r),其中 alpha = theta / (2 * (1 - theta))
  3. 对向量进行 L2 归一化
  4. 构建核矩阵 L = Diag(scores) * S * Diag(scores),其中 S[i][j] = (1 + dot(emb[i], emb[j])) / 2
  5. 运行贪心 DPP MAP 推理算法,选择多样性子集

注意事项

  • theta 必须在 [0, 1) 范围内
  • maxLength 必须为正整数
  • 向量列中的所有向量维度必须一致
  • 得分或向量为 NULL 的行会被自动跳过

rule_diversity

基于规则的多样性函数,通过贪心算法根据用户定义的规则对推荐结果进行打散重排,支持灵活的多样性约束配置。

函数签名

java
public CacheTable evaluate(
    CacheTable targetTable,
    CacheTable ruleTable,
    String maxReturn
)

参数说明

参数类型说明
targetTableCacheTable待打散的目标表
ruleTableCacheTable规则表,定义多样性约束规则
maxReturnString最大返回记录数

规则表字段说明

字段名类型说明
window_sizeInteger窗口大小
window_startInteger窗口起始位置(从 1 开始)
window_numInteger滑动窗口数量(1 表示不滑动)
diversity_columnString目标表中用于打散的列名
diversity_valueString匹配值(为空时约束适用于每个不同的值)
opString比较运算符(>=<
diversity_numInteger约束阈值
weightDouble规则权重,权重越高优先级越高

返回值:返回打散后的 CacheTable,结构与目标表相同。

使用示例

sql
-- 创建规则表
CACHE TABLE diversity_rules AS
SELECT
    5 AS window_size,
    1 AS window_start,
    1 AS window_num,
    'category' AS diversity_column,
    '' AS diversity_value,
    '<' AS op,
    2 AS diversity_num,
    1.0 AS weight
UNION ALL
SELECT
    3, 1, 1, 'brand', 'Nike', '=', 1, 2.0;

-- 基于规则打散,返回 20 条
CACHE TABLE rule_diversify_result AS
CALL rule_diversity(rec_item, diversity_rules, '20');

工作原理

  1. 解析规则表,为每条规则构建滑动窗口
  2. 对每个输出位置,贪心选择违反约束惩罚最小的未分配物品
  3. 如果没有违反约束,则按原始排名顺序选择
  4. 权重越高的规则,违反时的惩罚越大

注意事项

  • 规则表必须包含所有必需字段
  • diversity_column 必须在目标表中存在
  • diversity_value 为空时,约束适用于窗口内每个不同的属性值
  • 目标表中用于打散的列可以是单值或列表

json_to_table

JSON 转表函数,将 JSON 字符串转换为 CacheTable 表。

函数签名

java
public CacheTable evaluate(String jsonString)

参数说明

参数类型说明
jsonStringStringJSON 字符串,支持 JSON 对象或 JSON 数组

返回值:返回转换后的 CacheTable,列名和类型根据 JSON 内容自动推断。

使用示例

sql
-- 将 JSON 数组转换为表
CACHE TABLE json_result AS
CALL json_to_table('[{"id": 1, "name": "Alice"}, {"id": 2, "name": "Bob"}]');

-- 将单个 JSON 对象转换为表
CACHE TABLE single_obj AS
CALL json_to_table('{"id": 1, "name": "Alice", "score": 95.5}');

工作原理

  1. 解析 JSON 字符串,支持 JSON 对象和 JSON 数组
  2. 收集所有键作为列名,保持插入顺序
  3. 根据第一个非空值自动推断列类型(BOOLEANDOUBLEVARCHARARRAY<...>
  4. 将每个 JSON 对象转换为一行记录

注意事项

  • JSON 字符串不能为空
  • 必须是 JSON 对象或 JSON 数组格式
  • 数组中的嵌套对象会以 JSON 字符串形式存储为 VARCHAR
  • 数组类型会自动推断元素类型(ARRAY<DOUBLE>ARRAY<BOOLEAN>ARRAY<VARCHAR>

tag_to_vec

标签转向量函数,将标签列转换为 Multi-Hot 向量表示。

函数签名

java
public CacheTable evaluate(CacheTable input, String tagColName, String outputColName)

参数说明

参数类型说明
inputCacheTable输入表
tagColNameString标签列名,可以是单值或列表
outputColNameString输出向量列名

返回值:返回添加了向量列的 CacheTable,新列类型为 ARRAY<FLOAT>

使用示例

sql
-- 将用户标签转换为 Multi-Hot 向量
CACHE TABLE user_with_vec AS
CALL tag_to_vec(user_info, 'tags', 'tag_vector');

-- 将物品类目转换为向量
CACHE TABLE item_with_vec AS
CALL tag_to_vec(item_info, 'categories', 'category_vector');

工作原理

  1. 遍历所有行,收集标签列中的所有唯一标签,构建标签到索引的映射
  2. 为每行生成 Multi-Hot 向量,向量维度等于唯一标签数
  3. 如果行中包含某个标签,对应位置为 1.0,否则为 0.0
  4. 将向量列追加到原始表中

注意事项

  • 标签列可以是单值(字符串)或列表(ARRAY<STRING>
  • 输出列名不能与已有列名重复
  • 向量维度取决于所有行中唯一标签的总数

weighted_merge

加权合并函数,按指定权重将多个表合并为一个表,支持按主键去重。

函数签名

java
public CacheTable evaluate(String primaryKey, String weights, String limit, CacheTable... tables)

参数说明

参数类型说明
primaryKeyString主键列名,用于去重
weightsString各表权重,逗号分隔,如 "2,1,1"
limitString最大返回记录数
tablesCacheTable...一个或多个输入表,所有表结构必须相同

返回值:返回合并后的 CacheTable,结构与输入表相同。

使用示例

sql
-- 按权重 2:1:1 合并三个召回通道,返回 100 条
CACHE TABLE merged_recall AS
CALL weighted_merge('item_id', '2,1,1', '100', recall_channel_a, recall_channel_b, recall_channel_c);

-- 按权重 3:2 合并两个召回通道,返回 50 条
CACHE TABLE merged_result AS
CALL weighted_merge('item_id', '3,2', '50', recall_a, recall_b);

工作原理

  1. 使用 MergeUtils 的公共加权轮询内核,每轮按表的顺序,从每个表中取权重数量的有效记录
  2. 按主键去重,已出现的记录不再重复添加,且不占用该表本轮的权重配额
  3. 重复轮次直到达到 limit 或所有表遍历完毕
  4. 权重越大的表,每轮贡献的记录越多

注意事项

  • 所有输入表的结构(列名和类型)必须完全相同
  • primaryKey 为空字符串时不去重;指定主键时按该列的字符串值去重
  • 权重数量必须与表的数量一致
  • 权重和 limit 必须为正整数
  • 主键列必须存在于所有表中
  • weighted_merge 实现了 UnionLikeTableFunction,因此被视为类 UNION 合并节点。在 SQL 函数中,当 IGNORE_UNION_EXCEPTION=true 且某个输入缓存表的所有消费路径最终都进入 UNION/weighted_merge 时,该输入失败可降级为空表,其余输入继续合并;如果该输入还流向不经过合并的消费路径,则不会降级
  • 上述识别依赖编译期静态绑定的函数实例;通过 GET() 动态解析函数名的调用没有静态绑定实例,不会在编译期被识别为类 UNION 节点

call_sqlrec_api

远程 SQLRec API 调用函数,用于调用远端 SQLRec 实例上已发布的 API(即通过 CREATE API 暴露的 SQL 函数),并将返回结果转换为 CacheTable。适用于跨集群/跨实例调用其他 SQLRec 服务的场景。

函数签名

java
public CacheTable evaluate(ReadonlyContext context, String url, CacheTable... tables)

参数说明

参数类型说明
contextReadonlyContext只读上下文(自动注入),用于传递变量和指标标签
urlString远端 SQLRec API 地址,例如 http://host:port/api/v1/function_name
tablesCacheTable...一个或多个输入表,表名作为远端函数的输入占位符

返回值:返回包含远端函数执行结果的 CacheTable,列名和类型根据响应数据自动推断。

使用示例

sql
-- 准备输入数据
CACHE TABLE user_input AS
SELECT 1001 AS user_id, 'Alice' AS user_name;

-- 调用远端 SQLRec 实例上已发布的推荐 API
CACHE TABLE remote_rec AS
CALL call_sqlrec_api(
    'http://remote-sqlrec:30001/api/v1/recommend',
    user_input
);

SELECT * FROM remote_rec;

工作原理

  1. 将每个输入表按表名序列化为 Map<String, List<Map<String, Object>>>,连同上下文变量一起发送到远端 API
  2. 远端 SQLRec 执行对应的 SQL 函数并返回结果数据
  3. 根据响应行推断输出字段类型(BOOLEANDOUBLEVARCHARARRAY<...> 等)
  4. 复杂值(Map 等)会被序列化为 JSON 字符串存储为 VARCHARARRAY 类型保持为列表对象

注意事项

  • url 不能为空,且必须指向有效的 SQLRec API 端点
  • 至少需要传入一个输入表,且每个输入表必须有表名
  • 远端 API 调用失败(返回空数据或错误消息)时会抛出异常
  • 输入表的表名需与远端函数定义中的输入表占位符匹配

truncate_table

表截取函数,从输入表中截取指定范围的行记录。

函数签名

java
public CacheTable evaluate(CacheTable input, String start, String end)

参数说明

参数类型说明
inputCacheTable输入表
startString起始行索引(从 0 开始,包含)
endString结束行索引(不包含)

返回值:返回截取后的 CacheTable,结构与输入表相同。

使用示例

sql
-- 获取第 10 到 20 条记录
CACHE TABLE partial_result AS
CALL truncate_table(recall_item, '10', '20');

-- 获取前 100 条记录
CACHE TABLE top_100 AS
CALL truncate_table(recall_item, '0', '100');

注意事项

  • startend 必须为有效的整数字符串
  • startend 必须为非负数
  • start 必须小于或等于 end
  • 截取范围为左闭右开区间 [start, end)

get_variables

获取变量函数,从执行上下文中获取所有变量,返回一个包含变量键值对的表。

函数签名

java
public CacheTable evaluate(ReadonlyContext context)

参数说明

参数类型说明
contextReadonlyContext只读上下文

返回值:返回一个 2 列的 CacheTable,列名为 keyvalue,类型均为 VARCHAR

使用示例

sql
-- 设置一些变量
SET 'user_id' = '12345';
SET 'limit' = '100';

-- 获取所有变量
CACHE TABLE all_vars AS
CALL get_variables();

-- 查看变量
SELECT * FROM all_vars;

工作原理

  1. 从执行上下文中获取所有变量
  2. 将每个变量的键值对转换为一行记录
  3. 返回包含所有变量的表

set_variables

设置变量函数,从表中读取键值对并设置到执行上下文中。

函数签名

java
public CacheTable evaluate(ExecuteContext context, CacheTable input)

参数说明

参数类型说明
contextExecuteContext执行上下文
inputCacheTable输入表,必须恰好有 2 列,且均为字符串类型

返回值:返回输入表本身。

使用示例

sql
-- 创建变量表
CACHE TABLE var_table AS
SELECT 'user_id' AS key, '12345' AS value
UNION ALL
SELECT 'limit', '100';

-- 设置变量
CALL set_variables(var_table);

-- 使用设置的变量
SELECT `get`('user_id') AS user_id;

注意事项

  • 输入表必须恰好有 2 列
  • 两列都必须是字符串类型(VARCHAR 或 CHAR)
  • 第一列为变量名,第二列为变量值
  • 如果变量值为 NULL,则会删除该变量

feature_coverage_metrics

特征覆盖率打点函数,计算表中各字段的特征覆盖率并上报指标。

函数签名

java
public Void evaluate(ReadonlyContext context, String metricsName, CacheTable... tables)

参数说明

参数类型说明
contextReadonlyContext只读上下文
metricsNameString指标名称
tablesCacheTable...一个或多个输入表

返回值:无返回值。

使用示例

sql
-- 计算并上报特征覆盖率
CALL feature_coverage_metrics('feature.coverage', user_features, item_features);

-- 仅计算单个表的覆盖率
CALL feature_coverage_metrics('user.feature.coverage', user_info);

工作原理

  1. 遍历每个表的每个字段
  2. 统计每个字段的非空值数量(null、空 Collection、空 Map 视为缺失)
  3. 计算覆盖率 = 非空值数量 / 总行数
  4. 使用 summary 类型上报指标,tags 包含 table(表名)和 field(字段名)

注意事项

  • 如果表为空,则跳过该表
  • 指标名称不能为空

get_growthbook_features

GrowthBook 特征获取函数,从 GrowthBook 平台获取 A/B 实验特征值,并将实验参数设置为执行上下文变量,同时返回实验追踪数据用于指标计算。

函数签名

java
public CacheTable evaluate(ExecuteContext context, String apiHost, String clientKey,
                           CacheTable usertable, String... featureKeys)

参数说明

参数类型说明
contextExecuteContext执行上下文
apiHostStringGrowthBook API 地址
clientKeyStringGrowthBook 客户端密钥
usertableCacheTable用户表,表中的列将作为用户属性传入 GrowthBook
featureKeysString...一个或多个特征键名

返回值:返回包含实验追踪数据的 CacheTable,包含以下字段:

字段名类型说明
experiment_idVARCHAR实验标识
variation_idVARCHAR实验分组标识
user_idVARCHAR用户标识

使用示例

sql
-- 获取 GrowthBook 特征并设置变量
CACHE TABLE gb_tracking AS
CALL get_growthbook_features(
    'https://cdn.growthbook.io',
    'sdk-abc123',
    user_info,
    'new_recommendation_algo',
    'ui_theme'
);

-- 使用设置的实验变量
SELECT `get`('new_recommendation_algo') AS algo;

工作原理

  1. 根据 apiHostclientKey 创建或复用 GrowthBookClient(客户端会被缓存)
  2. 遍历用户表中的每一行,将行数据序列化为 JSON 作为用户属性
  3. 对每个特征键调用 evalFeature 获取特征值
  4. 如果特征有实验结果,将实验值通过 context.setVariable() 设置为变量,变量名即特征键名
  5. 收集实验追踪数据(实验 ID、分组 ID、用户 ID)并返回

注意事项

  • apiHostclientKey 不能为空
  • usertable 不能为空
  • 至少需要指定一个 featureKey
  • GrowthBookClient 初始化失败时会抛出异常
  • 同一组 apiHostclientKey 会复用同一个客户端实例

sleep

睡眠函数,使当前线程休眠指定的毫秒数。主要用于测试、限流或模拟延迟等场景。

函数签名

java
public Void evaluate(String millisStr)

参数说明

参数类型说明
millisStrString休眠时长,单位为毫秒,必须是非负整数字符串

返回值:无返回值。

使用示例

sql
-- 休眠 1000 毫秒(1 秒)
CALL sleep('1000');

-- 休眠 500 毫秒
CALL sleep('500');

工作原理

  1. 解析 millisStr 为长整型
  2. 调用 Thread.sleep() 使当前线程休眠指定时长
  3. 如果线程在休眠期间被中断,会恢复中断状态并抛出异常

注意事项

  • millisStr 必须是有效的长整数字符串
  • 休眠时长必须为非负数
  • 该函数仅产生副作用,无返回值

标量函数(Scalar Function)

array_contains

数组包含函数,检查数组是否包含指定元素。

函数签名

java
public static Boolean evaluate(List<?> list, Object element)

参数说明

参数类型说明
listList<?>输入数组
elementObject要检查的元素

返回值:如果数组包含该元素返回 true,否则返回 false;如果任一参数为 null 则返回 null

使用示例

sql
-- 检查用户标签是否包含 'vip'
SELECT
    user_id,
    array_contains(tags, 'vip') AS is_vip
FROM user_info;

-- 筛选包含特定标签的用户
SELECT *
FROM user_info
WHERE array_contains(tags, 'active') = true;

array_contains_all

数组全包含函数,检查数组是否包含所有指定元素。

函数签名

java
public static Boolean evaluate(List<?> list, List<?> elements)

参数说明

参数类型说明
listList<?>输入数组
elementsList<?>要检查的元素列表

返回值:如果数组包含所有指定元素返回 true,否则返回 false;如果任一参数为 null 则返回 null

使用示例

sql
-- 检查用户是否同时拥有多个标签
SELECT
    user_id,
    array_contains_all(tags, ARRAY['vip', 'active']) AS is_vip_active
FROM user_info;

-- 筛选同时满足多个条件的用户
SELECT *
FROM user_info
WHERE array_contains_all(tags, ARRAY['premium', 'verified']) = true;

array_contains_any

数组任一包含函数,检查数组是否包含指定元素中的任意一个。

函数签名

java
public static Boolean evaluate(List<?> list, List<?> elements)

参数说明

参数类型说明
listList<?>输入数组
elementsList<?>要检查的元素列表

返回值:如果数组包含任一指定元素返回 true,否则返回 false;如果任一参数为 null 则返回 null

使用示例

sql
-- 检查用户是否拥有任意一个 VIP 等级
SELECT
    user_id,
    array_contains_any(levels, ARRAY['gold', 'platinum', 'diamond']) AS is_high_level
FROM user_info;

-- 筛选拥有任意指定标签的用户
SELECT *
FROM user_info
WHERE array_contains_any(tags, ARRAY['new_user', 'trial']) = true;

random_vec

随机向量生成函数,生成指定维度的归一化随机向量。

函数签名

java
public List<Double> evaluate(String dimensionStr)

参数说明

参数类型说明
dimensionStrString向量维度,必须是正整数字符串

返回值:返回归一化的随机向量(List<Double>),向量的 L2 范数为 1。

使用示例

sql
-- 生成 64 维随机向量
SELECT
    user_id,
    random_vec('64') AS random_embedding
FROM user_info;

-- 为冷启动用户生成随机向量
CACHE TABLE cold_start_users AS
SELECT
    user_id,
    random_vec('128') AS user_embedding
FROM new_users;

工作原理

  1. 解析维度参数为整数
  2. 生成指定维度的随机向量
  3. 对向量进行 L2 归一化,使范数为 1

注意事项

  • 维度必须是正整数
  • 生成的向量已归一化,可直接用于相似度计算

uuid

UUID 生成函数,生成一个随机的 UUID 字符串。

函数签名

java
public String evaluate()

返回值:返回一个随机 UUID 字符串,格式如 ee073e63-b74a-4c7e-8fea-60459729099c

使用示例

sql
-- 生成请求 ID
CACHE TABLE request_meta AS
SELECT
    user_id,
    CAST(CURRENT_TIMESTAMP AS BIGINT) AS req_time,
    uuid() AS req_id
FROM user_info;

l2_norm

L2 归一化函数,对向量进行 L2 归一化处理。

函数签名

java
public List<Double> evaluate(Object vector)

参数说明

参数类型说明
vectorObject输入向量,必须是数字列表

返回值:返回归一化后的向量(List<Double>),使得向量的 L2 范数为 1。

使用示例

sql
-- 对用户向量进行归一化
CACHE TABLE normalized_user AS
SELECT
    user_id,
    l2_norm(user_embedding) AS normalized_embedding
FROM user_features;

工作原理

  1. 计算向量的 L2 范数:norm = sqrt(sum(x_i^2))
  2. 对每个元素除以范数:x_i' = x_i / norm
  3. 归一化后的向量常用于余弦相似度计算

ip

内积(Inner Product)计算函数,计算两个向量的内积(点积)。

函数签名

java
public Double evaluate(Object emb1, Object emb2)

参数说明

参数类型说明
emb1Object第一个向量,必须是数字列表
emb2Object第二个向量,必须是数字列表

返回值:返回两个向量的内积(Double)。

使用示例

sql
-- 计算用户向量和物品向量的内积
SELECT
    user_id,
    item_id,
    ip(user_embedding, item_embedding) AS similarity
FROM user_item_pairs;

-- 向量召回:按内积排序
CACHE TABLE vector_recall AS
SELECT item_embedding.id AS item_id
FROM user_embedding JOIN item_embedding ON 1=1
ORDER BY ip(user_embedding.embedding, item_embedding.embedding) DESC
LIMIT 300;

工作原理

  • 内积计算:ip = sum(emb1[i] * emb2[i])
  • 如果向量已归一化,内积等于余弦相似度
  • 常用于向量检索和相似度计算

get

变量获取函数,从执行上下文中获取变量的值。常用于在SQL中引用通过 SET 语句设置的变量。

函数签名

java
public static String evaluate(DataContext context, String key)

参数说明

参数类型说明
keyString变量名

返回值:返回变量的值(String),如果变量不存在则返回 NULL

注意

由于 get 是 SQL 关键字,使用时需要用反引号包裹函数名,写作 `get`

使用示例

sql
-- 设置变量
SET 'user_id' = '12345';

-- 获取变量值
SELECT `get`('user_id') AS user_id;

-- 在表达式中使用
SELECT `get`('user_id') || '_suffix' AS user_id_with_suffix;

-- 类型转换
SELECT CAST(`get`('limit_count') AS INT) AS limit_count;

-- 从表中获取变量名并使用
CACHE TABLE var_names AS SELECT 'user_id' AS var_name;
SELECT `get`(var_name) AS var_value FROM var_names;

工作原理

  1. 函数接收一个变量名作为参数
  2. 从执行上下文(ExecuteContext)中查找对应的变量值
  3. 返回变量值,如果变量不存在则返回 NULL

典型应用场景

  • 参数化SQL查询
  • 动态配置传递
  • 跨语句共享变量

get_or_default

变量获取函数(带默认值),从执行上下文中获取变量的值,如果变量不存在则返回指定的默认值。

函数签名

java
public static String evaluate(DataContext context, String key, String defaultValue)

参数说明

参数类型说明
keyString变量名
defaultValueString默认值,当变量不存在时返回

返回值:返回变量的值(String),如果变量不存在则返回 defaultValue

使用示例

sql
-- 设置变量
SET 'func_name' = 'add_col';

-- 获取变量值,如果不存在则使用默认值
SELECT `get_or_default`('user_id', 'default_user') AS user_id;

-- 动态调用函数:变量存在时使用变量值
CALL `get_or_default`('func_name', 'shuffle')(my_table);

-- 动态调用函数:变量不存在时使用默认值
CALL `get_or_default`('unknown_func', 'shuffle')(my_table);

工作原理

  1. 函数接收变量名和默认值两个参数
  2. 从执行上下文(ExecuteContext)中查找对应的变量值
  3. 如果变量存在,返回变量值;如果变量不存在,返回默认值

典型应用场景

  • 动态函数调用,提供兜底函数
  • 配置项获取,提供默认配置
  • 参数化 SQL,提供默认参数

自定义 UDF

可以参考 编程模型 文档了解如何开发自定义 UDF。