本文为您介绍 RecTemplate 项目中自定义的 ODPS 函数(UDF/UDAF/UDTF)的使用方法,包括函数简介、入参说明和输出说明。
UDF 函数汇总
UDF 函数(标量函数)
序号 | 函数名 | 函数类 | 功能简介 |
1 |
|
| 将数值字段分段,返回该数值所在区间的索引 |
2 |
|
| 将多个标签或类别字段进行组合,生成笛卡尔积形式的组合标签 |
3 |
|
| 根据 IP 地址解析地理位置信息(国家、省份、城市) |
4 |
|
| 计算纯数值格式的 CTR,支持 Wilson 置信区间校正 |
5 |
|
| 计算键值对格式的 CTR,支持 Wilson 置信区间校正 |
6 |
|
| 计算 MAP 类型输入的 CTR,功能同 CTRKV |
7 |
|
| 将 MAP 类型转换为 KV 格式的字符串 |
8 |
|
| 单序列补充功能,支持去重和属性处理 |
9 |
|
| 用历史序列补充实时序列,支持去重和时间戳处理 |
10 |
|
| 合并两个序列特征,限制最终长度 |
UDAF 函数(聚合函数)
序号 | 函数名 | 函数类 | 功能简介 |
1 |
|
| 对类别特征进行聚合统计,支持 sum/max/min/avg 等聚合方法 |
2 |
|
| 对类别进行键值统计聚合,返回 Top-K 结果(MAP 类型) |
3 |
|
| 按事件类型统计类别特征的键值对,支持 TopK 过滤 |
4 |
|
| 按时间戳排序生成最近的序列特征(TopK) |
UDTF 函数(表值函数)
序号 | 函数名 | 函数类 | 功能简介 |
1 |
|
| 根据用户行为日志实时生成序列特征 |
2 |
|
| 与 RtSeqFeature 类似,但去重逻辑略有差异 |
3 |
|
| 根据事件时间和滑动窗口配置,计算消息的延迟时间 |
4 |
|
| 基于滑动时间窗口统计用户的实时行为特征 |
5 |
|
| 将序列数组拆分为多行,每行包含序列长度倒序索引、元素值和列名 |
一、UDF 函数(标量函数)
1. Bound_Index
函数类:bound_index.BoundIndex
功能简介:将数值字段分段,返回该数值所在区间的索引。常用于数值特征的离散化处理。
入参说明
参数名 | 类型 | 必填 | 说明 |
field | DOUBLE | 是 | 需要分段的数值字段 |
boundaries | ARRAY | 是 | 分段点数组,例如 |
输出说明
输出 | 类型 | 说明 |
index | BIGINT | 数值所在区间的索引。如果 field < boundaries[0],返回 0;如果 boundaries[i] <= field < boundaries[i+1],返回 i+1;如果 field >= 最大值,返回 len(boundaries)。如果输入为 NULL,返回 -1 |
使用示例
-- 注册函数
CREATE FUNCTION bound_index AS 'bound_index.BoundIndex' USING 'bound_index.py';
-- 使用示例:将分数分段 [0, 60), [60, 70), [70, 80), [80, 90), [90, +∞)
SELECT bound_index(score, ARRAY(60.0, 70.0, 80.0, 90.0)) AS grade_level
FROM student_table;
-- score=55 返回 0, score=65 返回 1, score=85 返回 3, score=95 返回 42. Tag_Combo
函数类:tag_combo.TagCombo
功能简介:将多个标签或类别字段进行组合,生成笛卡尔积形式的组合标签。
入参说明
参数名 | 类型 | 必填 | 说明 |
tag_fields... | 任意类型(可变参数) | 是 | 多个标签或类别字段,每个字段可以包含多个标签(使用 |
输出说明
输出 | 类型 | 说明 |
combo_tags | STRING | 组合标签字符串,使用 |
使用示例
-- 注册函数
CREATE FUNCTION tag_combo AS 'tag_combo.TagCombo' USING 'tag_combo.py';
-- 组合多个标签
SELECT tag_combo(category, brand, style) AS combo_result
FROM item_table;
-- 输入: category="服装\x1d鞋帽", brand="Nike\x1dAdidas"
-- 输出: "服装_Nike\x1d服装_Adidas\x1d鞋帽_Nike\x1d鞋帽_Adidas"3. AnalysisIP - IP 地址解析函数
函数类:com.aliyun.pai.udf.AnalysisIP
功能简介:根据输入的 IP 地址,解析其对应的地理位置信息(国家、省份、城市)。使用 IPIP 离线 IP 数据库 (ipipfree.ipdb) 进行解析。
入参说明
参数名 | 类型 | 必填 | 说明 |
| String | 是 | 需要解析的 IP 地址,如 |
| String | 是 | 需要返回的地理位置字段,可选值: - |
输出说明
返回类型 | 说明 |
String | 对应 location 字段的地理位置名称。如果解析失败或字段不存在,返回 |
使用示例
SELECT AnalysisIP('114.114.114.114', 'country') AS country,
AnalysisIP('114.114.114.114', 'region') AS province,
AnalysisIP('114.114.114.114', 'city') AS city
FROM dual;
-- 返回: 中国, 江苏省, 南京市注意事项
需要在 ODPS 中注册资源文件
ipipfree.ipdb。
4. CTR_NUM
函数类:ctr_num.CTRNum
功能简介:计算纯数值格式的 CTR,支持 Wilson 置信区间校正。
入参说明
参数名 | 类型 | 必填 | 说明 |
clicks | DOUBLE | 是 | 点击数 |
impressions | DOUBLE | 是 | 展示数 |
is_wilson | BOOLEAN | 否 | 是否使用 Wilson 置信区间,默认 TRUE,置信区间0.95 |
输出说明
输出 | 类型 | 说明 |
ctr | DOUBLE | CTR 值,范围 [0, 1] |
使用示例
-- 注册函数
CREATE FUNCTION ctr_num AS 'ctr_num.CTRNum' USING 'ctr_num.py';
-- 计算各物品的 CTR
SELECT item_id,
ctr_num(click_count, impression_count, true) AS ctr
FROM item_stats_table
ORDER BY ctr DESC;5. CTR_KV
函数类:ctr_kv.CTRKV
功能简介:计算键值对格式的 CTR(点击通过率),支持 Wilson 置信区间校正。适用于小样本场景的稳健 CTR 估算。
入参说明
参数名 | 类型 | 必填 | 说明 |
clicks | STRING | 是 | 点击数,格式为 KV 字符串,例如 |
impressions | STRING | 是 | 展示数,格式同 clicks |
is_wilson | BOOLEAN | 否 | 是否使用 Wilson 置信区间,默认 TRUE |
输出说明
输出 | 类型 | 说明 |
ctr_kv | STRING | CTR 键值对字符串,按 CTR 值降序排列,格式: |
使用示例
-- 注册函数
CREATE FUNCTION ctr_kv AS 'ctr_kv.CTRKV' USING 'ctr_kv.py';
-- 计算各物品的 CTR
SELECT ctr_kv(clicks_str, impressions_str, true) AS ctr_result
FROM item_stats_table;
-- 输出示例: item1:0.152341\x1ditem2:0.098234\x1ditem3:0.0756126. CTR_KV_MAP
函数类:ctr_kv_map.CTRKVMap
功能简介:计算 MAP 类型输入的 CTR,功能同 CTRKV,但输入输出为 MAP 类型。
入参说明
参数名 | 类型 | 必填 | 说明 |
clicks | MAP<STRING, FLOAT> | 是 | 点击数 MAP |
impressions | MAP<STRING, FLOAT> | 是 | 展示数 MAP |
is_wilson | BOOLEAN | 否 | 是否使用 Wilson 置信区间,默认 TRUE |
输出说明
输出 | 类型 | 说明 |
ctr_map | MAP<STRING, FLOAT> | CTR MAP,按 CTR 值降序排列 |
使用示例
-- 注册函数
CREATE FUNCTION ctr_kv_map AS 'ctr_kv_map.CTRKVMap' USING 'ctr_kv_map.py';
-- 计算各物品的 CTR(MAP 版本)
SELECT ctr_kv_map(clicks_map, impressions_map, true) AS ctr_map
FROM item_stats_table;7. Map2Str
函数类:map2str.Map2Str
功能简介:将 MAP 类型转换为 KV 格式的字符串。
入参说明
参数名 | 类型 | 必填 | 说明 |
field | MAP<STRING, FLOAT> | 是 | 需要转换的 MAP 字段 |
输出说明
输出 | 类型 | 说明 |
kv_str | STRING | KV 格式字符串,使用 |
使用示例
-- 注册函数
CREATE FUNCTION map2str AS 'map2str.Map2Str' USING 'map2str.py';
-- 将 MAP 转换为字符串
SELECT map2str(feature_map) AS feature_str
FROM feature_table;
-- 输入: {'age': 25.0, 'gender': 1.0}
-- 输出: age:25.0\x1dgender:1.08. SINGLE_SEQ_REPLENISH
函数类:single_seq_replenish.SingleSeqReplenish
功能简介:单序列补充功能,支持更复杂的去重和属性处理逻辑。
入参说明
参数名 | 类型 | 必填 | 说明 |
seq1 | STRING | 是 | 实时序列特征 |
re_seq | STRING | 是 | T-1 序列特征 |
id1 | STRING | 是 | 重复序列 ID1(用于去重判断) |
id2 | STRING | 是 | 重复序列 ID2(用于去重判断) |
curren_time | 任意类型 | 否 | 当前事件时间戳,默认 0 |
delim | STRING | 否 | 分隔符,默认 |
length | BIGINT | 否 | 最终序列长度,默认 50 |
deduplication_method | BIGINT | 否 | 去重方法:1=按 item+event 去重,2=按 item+event+event_time 去重,默认 1 |
输出说明
输出 | 类型 | 说明 |
replenished_seq | STRING | 补充后的序列字符串 |
使用示例
-- 注册函数
CREATE FUNCTION single_seq_replenish AS 'single_seq_replenish.SingleSeqReplenish' USING 'single_seq_replenish.py';
-- 单序列补充
SELECT user_id,
single_seq_replenish(seq1, seq2, id1, id2, event_time, ';', 50, 1) AS full_seq
FROM user_sequence_table;9. SEQ_REPLENISH
函数类:seq_replenish.SeqReplenish
功能简介:用历史序列补充实时序列,支持去重和时间戳处理。用于构建完整的用户行为序列。tf流程多使用该函数
入参说明
参数名 | 类型 | 必填 | 说明 |
seq1 | STRING | 是 | 实时序列特征 |
re_seq | STRING | 是 | T-1 序列特征(需要补充的历史序列) |
curren_time | 任意类型 | 是 | 当前事件时间戳(Unix 时间) |
delim | STRING | 否 | 分隔符,默认 |
length | BIGINT | 否 | 最终序列长度,默认 50 |
deduplication_method | BIGINT | 否 | 去重方法:1=按 item+event 去重,2=按 item+event+event_time 去重,默认 1 |
输出说明
输出 | 类型 | 说明 |
replenished_seq | STRING | 补充后的序列字符串,已去重并按规则添加时间戳 |
使用示例
-- 注册函数
CREATE FUNCTION seq_replenish AS 'seq_replenish.SeqReplenish' USING 'seq_replenish.py';
-- 用历史序列补充实时序列
SELECT user_id,
seq_replenish(realtime_seq, history_seq, event_time, ';', 50, 1) AS full_seq
FROM user_sequence_table;10. SEQ_CONCAT
函数类:seq_concat.SeqConcat
功能简介:合并两个序列特征,将 seq2 补充到 seq1 后面,并限制最终长度。常用于实时特征和历史特征的拼接。tf的流程多使用该函数
入参说明
参数名 | 类型 | 必填 | 说明 |
seq1 | STRING | 是 | 实时序列特征 |
seq2 | STRING | 是 | T-1 序列特征(历史特征) |
delim | STRING | 否 | 分隔符,默认 |
length | BIGINT | 否 | 最终序列长度,默认 50 |
输出说明
输出 | 类型 | 说明 |
merged_seq | STRING | 合并后的序列,seq1 在前,seq2 在后,超过 length 则截断 |
使用示例
-- 注册函数
CREATE FUNCTION seq_concat AS 'seq_concat.SeqConcat' USING 'seq_concat.py';
-- 合并实时和历史行为序列
SELECT user_id,
seq_concat(realtime_seq, history_seq, ';', 100) AS full_seq
FROM user_sequence_table;
-- realtime_seq: "item1;item2;item3"
-- history_seq: "item4;item5;item6"
-- 输出: "item1;item2;item3;item4;item5;item6"二、UDAF 函数(聚合函数)
1. count_kv
函数类:count_kv.CountKV
功能简介:对类别特征进行聚合统计,支持 sum/max/min/avg 等聚合方法,并支持 TopK 过滤。
入参说明
参数名 | 类型 | 必填 | 说明 |
cate | STRING | 是 | 类别字段 |
value | DOUBLE | 是 | 需要累积的数值 |
agg_type | STRING | 是 | 聚合方法: |
top_k | BIGINT | 否 | 保留的最大键数,默认 50 |
is_apply_agg | BOOLEAN | 否 | 是否参与统计,默认 TRUE |
输出说明
输出 | 类型 | 说明 |
kv_str | STRING | 键值对字符串,按聚合值降序排列,使用 |
使用示例
-- 注册函数
CREATE FUNCTION count_kv AS 'count_kv.CountKV' USING 'count_kv.py';
-- 统计各类别的点击总数
SELECT category,
count_kv(item_id, click_count, 'sum', 100, true) AS item_stats
FROM item_click_table
GROUP BY category;
-- 输出示例: "item1:150\x1ditem2:120\x1ditem3:95"
-- 统计各类别的平均停留时长
SELECT category,
count_kv(item_id, duration, 'avg', 50, true) AS avg_duration_stats
FROM item_view_table
GROUP BY category;
-- 输出示例: "item1:45.5\x1ditem2:32.1\x1ditem3:28.7"2. count_kv_map
函数类:com.aliyun.pai.udf.CountKV
功能简介:对类别(Category)进行键值统计聚合,支持多种聚合方式(SUM/AVG/MAX/MIN),并返回 Top-K 结果。支持多类别同时聚合(类别之间用 \u001D 分隔)。
入参说明
参数名 | 类型 | 必填 | 说明 |
| String | 是 | 类别字段,支持多类别(用 |
| Float | 是 | 需要聚合的数值 |
| String | 否 | 聚合类型,可选值: - |
| Int | 否 | 保留 Top-K 个结果,默认 |
| Boolean | 是 | 是否应用聚合,为 |
输出说明
返回类型 | 说明 |
Map<String, Float> | 按聚合值降序排列的 Top-K 键值对 Map |
使用示例
-- 统计每个用户点击次数最多的 Top 10 商品类别
SELECT user_id,
CountKV(category, 1.0, 'sum', 10, true) AS top_categories
FROM user_behavior
GROUP BY user_id;
-- 返回示例: {"electronics": 15.0, "clothing": 10.0, "food": 8.0}注意事项
聚合类型为
avg时,会同时维护总和和计数,最终计算平均值。类别字段支持数组形式的输入,用特殊字符
\u001D(ASCII 29) 分隔。topK <= 0时返回所有类别## 2. CountKV - 类别键值聚合函数。
3. COUNT_CATES_KVS
函数类:count_cates_kvs.CountCatesKVS
功能简介:按事件类型统计类别特征的键值对,支持 TopK 过滤和事件计算公式。用于构建用户行为统计特征。
入参说明
参数名 | 类型 | 必填 | 说明 |
cates | 任意类型 | 是 | 类别字段,可以是单个值或多值(多值使用分隔符分隔) |
values | BIGINT/DOUBLE/ARRAY | 是 | 需要累积的数值,可以是单个值或列表 |
event | STRING | 是 | 事件字段 |
event_types | STRING/ARRAY | 是 | 需要统计的事件类型,多事件用 ` |
event_formulas | STRING | 否 | 事件计算公式,例如 |
sep | STRING | 否 | 多值类别的分隔符,默认 |
top_k | BIGINT | 否 | 保留的最大键数,默认 50 |
is_punish | BOOLEAN | 否 | 计算公式时是否用均值惩罚分母,默认 FALSE |
输出说明
输出 | 类型 | 说明 |
kv_array | ARRAY | 键值对数组,每个元素格式为 |
使用示例
-- 注册函数
CREATE FUNCTION count_cates_kvs AS 'count_cates_kvs.CountCatesKVS' USING 'count_cates_kvs.py';
-- 统计用户在不同事件下的类别行为
SELECT user_id,
count_cates_kvs(category, 1, event_type, 'click,purchase', 'purchase/click', '\x1d', 100, false) AS stats_array
FROM user_behavior_table
GROUP BY user_id;
-- 输出示例: ["cate1:10\x1dcate2:5", "cate1:2\x1dcate2:1"] (分别对应 click 和 purchase/click 的统计)4. WM_CONCAT_BY_SORT
函数类:wm_concat_by_sort.WMConcatBySort
功能简介:按时间戳排序生成最近的序列特征(TopK)。用于构建用户行为序列特征。
入参说明
参数名 | 类型 | 必填 | 说明 |
values | STRING | 是 | 需要连接的序列元素 |
timestamp | 任意类型 | 是 | 排序时间戳 |
event | STRING | 是 | 当前行事件 |
valid_event_selections | STRING/ARRAY | 是 | 有效事件选择列表,只有在此列表中的事件才会被计入序列 |
sequence_delim | STRING | 否 | 序列分隔符,默认 |
top_k | BIGINT | 否 | 序列长度,默认 50 |
输出说明
输出 | 类型 | 说明 |
seq_str | STRING | 按时间戳降序排列的前 TopK 个序列元素,使用 sequence_delim 分隔 |
使用示例
-- 注册函数
CREATE FUNCTION wm_concat_by_sort AS 'wm_concat_by_sort.WMConcatBySort' USING 'wm_concat_by_sort.py';
-- 构建用户最近的行为序列
SELECT user_id,
wm_concat_by_sort(item_id, event_time, event_type, 'click,view,purchase', ';', 50) AS behavior_seq
FROM user_behavior_table
GROUP BY user_id;
-- 输出示例: "item50;item49;item48;...;item1" (按时间降序,最近的行为在前)三、UDTF 函数(表值函数)
1. RtSeqFeature - 实时序列特征生成函数
函数类:com.aliyun.pai.udf.CountKV.RtSeqFeature
功能简介:根据用户行为日志实时生成序列特征。支持多事件类型、多场景、序列去重、长度控制、时间窗口防止数据穿越等功能。输出包含用户 ID、时间戳以及各个序列特征字段。
入参说明
参数序号 | 参数名 | 类型 | 必填 | 说明 |
1 |
| String | 是 | 用户 ID |
2 |
| Long (bigint) | 是 | 事件时间戳 |
3 |
| Array | 是 | 需要记录的序列特征字段值数组 |
4 |
| Array | 是 | 序列特征字段名称数组,如 |
5 |
| String | 是 | 当前事件类型 |
6 |
| Array | 是 | 有效事件类型列表,如 |
7 |
| Array | 是 | 每个事件对应的序列长度限制 |
8 |
| Int | 是 | 去重方法: - |
9 |
| String | 是 | 序列合并方法: - |
10 |
| Int | 是 | 防止数据穿越的秒数,如 |
11 |
| String | 是 | 序列分隔符,如 |
12 |
| String | 否 | 请求 ID,用于去重输出 |
13 |
| String | 否 | 当前场景值 |
14 |
| Array | 否 | 有效场景列表 |
参数说明:
当
args.length == 11: 无 requestid,无 scene当
args.length == 12: 有 requestid,无 scene当
args.length == 13: 无 requestid,有 scene当
args.length == 14: 有 requestid,有 scene
输出说明
返回类型 | 说明 |
Map<String, String> | 包含以下键值对: - |
输出序列键名格式:
如果
combine_method = "all":{all}_{max_seq_len}_seq__{field_name}如果
combine_method != "all":{event}_{seq_len}_seq__{field_name}如果有场景: 后缀加上
_{scene},如all_10_seq_home__item_id
使用示例
-- 生成用户实时行为序列特征
SELECT output.*
FROM user_behavior
LATERAL VIEW RtSeqFeature(
user_id, -- pid
event_timestamp, -- event_time
ARRAY(item_id, category), -- fields
ARRAY("item_id", "category"), -- field_names
event_type, -- event
ARRAY("click", "add", "buy"), -- valid_events
ARRAY(10, 8, 6), -- seq_len
1, -- duplicate_method
"all", -- combine_method
300, -- pre_seconds
";", -- sequence_delim
request_id, -- requestid (可选)
scene, -- scene (可选)
ARRAY("home", "search") -- valid_scenes (可选)
) output AS feature_map;输出示例
{
"pid": "user123",
"event_time": "1709280000",
"all_10_seq__item_id": "item5;item3;item1",
"all_10_seq__category": "electronics;clothing;food",
"all_10_seq__ts": "0;120;300"
}注意事项
该函数是 UDTF,会输出一行或多行数据。
支持按 requestid 去重输出,避免同一请求重复。
pre_seconds用于防止数据穿越,确保只使用历史数据。序列按时间倒序输出(最新的在前)。
2. RtSeqFeature2
函数类:com.aliyun.pai.udf.CountKV.RtSeqFeature2
功能简介:与 RtSeqFeature 功能基本相同,但是去重逻辑略有差异。RtSeqFeature2 在 duplicate_method=2 时,去重键为 fields[0] + event_time,而 RtSeqFeature 为 event + fields[0] + event_time。
入参说明
与 RtSeqFeature 完全相同。
输出说明
与 RtSeqFeature 完全相同。
与 RtSeqFeature 的区别
特性 | RtSeqFeature | RtSeqFeature2 |
去重键 (duplicate_method=2) |
|
|
去重键 (duplicate_method=1) |
|
|
使用场景:
如果需要根据 item 去重而不考虑事件类型,使用
RtSeqFeature2。如果需要同时考虑事件类型和 item,使用
RtSeqFeature。
3. MessageDelay
函数类:com.aliyun.pai.udf.CountKV.MessageDelay
功能简介:根据事件时间和滑动窗口配置,计算消息的延迟时间。用于实时特征生成时,为不同时间窗口的统计数据生成延迟行。
入参说明
参数序号 | 参数名 | 类型 | 必填 | 说明 |
1 |
| Array | 是 | 需要输出的基础字段数组 |
2 |
| Long (bigint) | 是 | 事件时间戳 |
3 |
| Array | 是 | 时间窗口列表 (秒),如 |
4 |
| Boolean | 是 | 是否保存 (生成延迟行): - |
输出说明
返回类型 | 说明 |
Array | 输出数组,包含: - 原始 fields 的所有字段 - 倒数第二个: 窗口触发时间戳 ( |
使用示例
-- 为不同时间窗口生成延迟数据行
SELECT output.*
FROM user_events
LATERAL VIEW MessageDelay(
ARRAY(user_id, event_type), -- fields
event_timestamp, -- event_time
ARRAY(1800, 3600, 7200), -- windows (30m, 1h, 2h)
true -- is_save
) output AS field1, field2, trigger_time, window_size;
-- 如果 is_save=true,会为每个窗口输出一行:
-- [user_id, event_type, 1709281800, 1800]
-- [user_id, event_type, 1709283600, 3600]
-- [user_id, event_type, 1709287200, 7200]注意事项
该函数是 UDTF,根据
is_save参数可能输出多行。窗口值会自动添加
0作为第一个窗口 (立即触发)。用于实时特征 MapReduce 的 Mapper 阶段,为 Reducer 生成延迟数据行。
4. SWCountCatesKVS
函数类:com.aliyun.pai.udf.CountKV.SWCountCatesKVS
功能简介:基于滑动时间窗口统计用户的实时行为特征,支持多时间窗口、多事件类型、多场景、Top-K 类别统计。用于实时特征批量生成场景。
入参说明
参数序号 | 参数名 | 类型 | 必填 | 说明 |
1 |
| String | 是 | 用户 ID |
2 |
| Array | 是 | 类别/标签字段数组,如 |
3 |
| Array | 是 | 类别字段名称数组,如 |
4 |
| Int | 是 | 统计值 (通常为 1) |
5 |
| String | 是 | 当前事件类型 |
6 |
| Array | 是 | 有效事件类型列表 |
7 |
| Int | 是 | 延迟标志: - |
8 |
| Array | 是 | 时间窗口列表 (秒),如 |
9 |
| Int | 是 | 每个类别保留的 Top-K 数量 |
10 |
| Int | 是 | 防止数据穿越的秒数 |
11 |
| Long (bigint) | 是 | 事件时间戳 |
12 |
| String | 是 | 输出日期 (格式: |
13 |
| String | 是 | 字段归属标识,如 |
14 |
| String | 否 | 请求 ID,用于去重输出 |
15 |
| String | 否 | 当前场景值 |
16 |
| Array | 否 | 有效场景列表 |
参数说明:
当
args.length == 13:无 requestid,无 scene。当
args.length == 14:有 requestid,无 scene。当
args.length == 15:无 requestid,有 scene。当
args.length == 16:有 requestid,有 scene。
输出说明
返回类型 | 说明 |
Map<String, String> | 包含以下键值对: - |
输出键名格式:
窗口名称格式:
rt{X}h{Y}m{Z}s,如rt30m、rt3h、rt24h。事件计数:
user__cnt_click_rt30m。类别键值对:
user__kv_style_id_click_rt30m。如果有场景:加上场景后缀,如
user__cnt_click_home_rt30m。
类别键值对格式:
使用
(char)29(即\u001D) 分隔多个键值对。每个键值对格式为
key:value。示例:
"cate1:15.0\u001Dcate2:10.0\u001Dcate3:8.0"。
使用示例
-- 生成用户实时行为统计特征
SELECT output.*
FROM user_behavior
LATERAL VIEW SWCountCatesKVS(
user_id, -- pid
ARRAY(style_id, goods_id), -- cates
ARRAY("style_id", "goods_id"), -- cate_names
1, -- value
event_type, -- event
ARRAY("click", "add", "buy"), -- valid_events
0, -- delay_flag (0表示立即生效)
ARRAY(1800, 3600, 7200, 86400), -- windows (30m, 1h, 2h, 24h)
100, -- topk
5, -- pre_seconds
event_timestamp, -- event_time
"20240101", -- output_date
"user", -- side
request_id, -- requestid (可选)
scene, -- scene (可选)
ARRAY("home", "search") -- valid_scenes (可选)
) output AS feature_map;输出示例
{
"pid": "user123",
"event_time": "1709280000",
"user__cnt_click_rt30m": "15",
"user__cnt_click_rt3h": "45",
"user__cnt_add_rt30m": "8",
"user__kv_style_id_click_rt30m": "style1:10.0\u001Dstyle2:5.0",
"user__kv_goods_id_click_rt30m": "goods1:8.0\u001Dgoods2:7.0"
}注意事项
该函数是 UDTF,输出包含统计结果。
delay_flag参数用于实现滑动窗口:0: 表示数据立即添加到所有窗口。非
0: 表示对应窗口的数据需要减去(用于维护滑动窗口)。
pre_seconds用于防止数据穿越。output_date用于控制输出起始日期,避免输出历史数据。支持按 requestid 去重输出。
类别值支持多标签,用
(char)29(即\u001D) 分隔。
5. SEQ_SPLIT
函数类:seq_split.SeqSplit
功能简介:将序列数组拆分为多行,每行包含序列长度倒序索引、元素值和列名。用于序列特征的展开处理。
入参说明
参数名 | 类型 | 必填 | 说明 |
seqs | ARRAY | 是 | 需要拆分的序列数组 |
col_names | ARRAY | 是 | 列名数组,与 seqs 一一对应 |
sep | STRING | 是 | 序列内部分隔符 |
输出说明
输出列 | 类型 | 说明 |
seq_len | BIGINT | 序列长度倒序索引(从后往前的位置) |
item | STRING | 序列元素值 |
col_name | STRING | 对应的列名 |
使用示例
-- 注册函数
CREATE FUNCTION seq_split AS 'seq_split.SeqSplit' USING 'seq_split.py';
-- 将序列拆分为多行
SELECT t.seq_len, t.item, t.col_name
FROM sequence_table
CROSS APPLY seq_split(ARRAY(seq1, seq2), ARRAY('seq1_col', 'seq2_col'), ';') AS t(seq_len, item, col_name);
-- 输入: seq1="item1;item2;item3", seq2="item4;item5"
-- 输出:
-- seq_len | item | col_name
-- --------|-------|----------
-- 3 | item1 | seq1_col
-- 2 | item2 | seq1_col
-- 1 | item3 | seq1_col
-- 2 | item4 | seq2_col
-- 1 | item5 | seq2_col四、函数注册和使用指南
1. 添加资源文件
在使用自定义函数之前,需要先将 Python 文件添加到 ODPS 资源中:
-- 添加单个 Python 文件
ADD py /path/to/bound_index.py;
-- 添加jar文件
ADD jar /path/to/feature-generate-mr-v1.21.jar;
-- 添加依赖的 JSON 配置文件(如 FG 配置)
ADD file /path/to/fg_config.json;2. 创建函数
-- 创建 UDF 函数
CREATE FUNCTION bound_index AS 'bound_index.BoundIndex' USING 'bound_index.py';
CREATE FUNCTION ctr_num AS 'ctr_num.CTRNum' USING 'ctr_num.py';
-- 创建 UDAF 函数
CREATE FUNCTION count_kv AS 'count_kv.CountKV' USING 'count_kv.py';
-- 创建 UDTF 函数
CREATE FUNCTION seq_split AS 'seq_split.SeqSplit' USING 'seq_split.py';3. 查看函数列表
-- 查看所有自定义函数
LIST FUNCTIONS;
-- 查看函数详细信息
DESCRIBE FUNCTION bound_index;4. 删除函数
-- 删除函数
DROP FUNCTION bound_index;5. 使用函数
-- 使用 UDF
SELECT bound_index(score, ARRAY(60.0, 70.0, 80.0, 90.0)) AS grade_level
FROM student_table;
-- 使用 UDAF
SELECT category, count_kv(item_id, click_count, 'sum', 100, true) AS item_stats
FROM item_click_table
GROUP BY category;
-- 使用 UDTF
SELECT t.seq_len, t.item, t.col_name
FROM sequence_table
CROSS APPLY seq_split(ARRAY(seq1), ARRAY('seq1_col'), ';') AS t(seq_len, item, col_name);五、常见问题
1. 分隔符说明
自定义函数中使用的主要分隔符:
\x1d(chr(29)): Group Separator,用于分隔键值对。;: 分号,用于分隔序列中的元素。|: 竖线,用于分隔时间序列或 FG 特征。#: 井号,用于序列分隔项和属性。
2. Wilson 置信区间
CTR 相关函数使用 Wilson 置信区间计算 CTR,公式为:
CTR = (p̂ + z²/(2n) - z * √(p̂(1-p̂)/n + z²/(4n²))) / (1 + z²/n)其中:
p̂ = clicks / impressions(观测点击率)n = impressions(展示数)z = 1.96(95% 置信水平)
优点:
对小样本情况更稳健。
避免极端值(0 或 1)的影响。
适合推荐系统中稀有物品的质量评估。