ODPS自定义函数

更新时间:
复制 MD 格式

本文为您介绍 RecTemplate 项目中自定义的 ODPS 函数(UDF/UDAF/UDTF)的使用方法,包括函数简介、入参说明和输出说明。

UDF 函数汇总

UDF 函数(标量函数)

序号

函数名

函数类

功能简介

1

Bound_Index

bound_index.BoundIndex

将数值字段分段,返回该数值所在区间的索引

2

Tag_Combo

tag_combo.TagCombo

将多个标签或类别字段进行组合,生成笛卡尔积形式的组合标签

3

AnalysisIP

com.aliyun.pai.udf.AnalysisIP

根据 IP 地址解析地理位置信息(国家、省份、城市)

4

CTR_NUM

ctr_num.CTRNum

计算纯数值格式的 CTR,支持 Wilson 置信区间校正

5

CTR_KV

ctr_kv.CTRKV

计算键值对格式的 CTR,支持 Wilson 置信区间校正

6

CTR_KV_MAP

ctr_kv_map.CTRKVMap

计算 MAP 类型输入的 CTR,功能同 CTRKV

7

Map2Str

map2str.Map2Str

将 MAP 类型转换为 KV 格式的字符串

8

SINGLE_SEQ_REPLENISH

single_seq_replenish.SingleSeqReplenish

单序列补充功能,支持去重和属性处理

9

SEQ_REPLENISH

seq_replenish.SeqReplenish

用历史序列补充实时序列,支持去重和时间戳处理

10

SEQ_CONCAT

seq_concat.SeqConcat

合并两个序列特征,限制最终长度

UDAF 函数(聚合函数)

序号

函数名

函数类

功能简介

1

count_kv

count_kv.CountKV

对类别特征进行聚合统计,支持 sum/max/min/avg 等聚合方法

2

count_kv_map

com.aliyun.pai.udf.CountKV

对类别进行键值统计聚合,返回 Top-K 结果(MAP 类型)

3

COUNT_CATES_KVS

count_cates_kvs.CountCatesKVS

按事件类型统计类别特征的键值对,支持 TopK 过滤

4

WM_CONCAT_BY_SORT

wm_concat_by_sort.WMConcatBySort

按时间戳排序生成最近的序列特征(TopK)

UDTF 函数(表值函数)

序号

函数名

函数类

功能简介

1

RtSeqFeature

com.aliyun.pai.udf.CountKV.RtSeqFeature

根据用户行为日志实时生成序列特征

2

RtSeqFeature2

com.aliyun.pai.udf.CountKV.RtSeqFeature2

与 RtSeqFeature 类似,但去重逻辑略有差异

3

MessageDelay

com.aliyun.pai.udf.CountKV.MessageDelay

根据事件时间和滑动窗口配置,计算消息的延迟时间

4

SWCountCatesKVS

com.aliyun.pai.udf.CountKV.SWCountCatesKVS

基于滑动时间窗口统计用户的实时行为特征

5

SEQ_SPLIT

seq_split.SeqSplit

将序列数组拆分为多行,每行包含序列长度倒序索引、元素值和列名

一、UDF 函数(标量函数)

1. Bound_Index

函数类bound_index.BoundIndex
功能简介:将数值字段分段,返回该数值所在区间的索引。常用于数值特征的离散化处理。

入参说明

参数名

类型

必填

说明

field

DOUBLE

需要分段的数值字段

boundaries

ARRAY

分段点数组,例如 [0.5, 0.7, 0.9]

输出说明

输出

类型

说明

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 返回 4

2. Tag_Combo

函数类tag_combo.TagCombo
功能简介:将多个标签或类别字段进行组合,生成笛卡尔积形式的组合标签。

入参说明

参数名

类型

必填

说明

tag_fields...

任意类型(可变参数)

多个标签或类别字段,每个字段可以包含多个标签(使用 \x1d 分隔)

输出说明

输出

类型

说明

combo_tags

STRING

组合标签字符串,使用 _ 连接不同标签,使用 \x1d 分隔组合结果

使用示例

-- 注册函数
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) 进行解析。

入参说明

参数名

类型

必填

说明

ip

String

需要解析的 IP 地址,如 "114.114.114.114"

location

String

需要返回的地理位置字段,可选值: - "country": 国家名称 - "region": 省份名称 - "city": 城市名称

输出说明

返回类型

说明

String

对应 location 字段的地理位置名称。如果解析失败或字段不存在,返回 null

使用示例

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 字符串,例如 item1:10\x1ditem2:5(使用 \x1d 分隔)

impressions

STRING

展示数,格式同 clicks

is_wilson

BOOLEAN

是否使用 Wilson 置信区间,默认 TRUE

输出说明

输出

类型

说明

ctr_kv

STRING

CTR 键值对字符串,按 CTR 值降序排列,格式:item1:0.15\x1ditem2:0.10

使用示例

-- 注册函数
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.075612

6. 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 格式字符串,使用 \x1d 分隔,格式:key1:value1\x1dkey2:value2

使用示例

-- 注册函数
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.0

8. 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

聚合方法:sum, max, min, avg

top_k

BIGINT

保留的最大键数,默认 50

is_apply_agg

BOOLEAN

是否参与统计,默认 TRUE

输出说明

输出

类型

说明

kv_str

STRING

键值对字符串,按聚合值降序排列,使用 \x1d 分隔

使用示例

-- 注册函数
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 分隔)。

入参说明

参数名

类型

必填

说明

cate

String

类别字段,支持多类别(用 \u001D 分隔),如 "cate1\u001Dcate2"

value

Float

需要聚合的数值

aggType

String

聚合类型,可选值: - "sum": 求和 (默认) - "avg": 平均值 - "max": 最大值 - "min": 最小值

topK

Int

保留 Top-K 个结果,默认 20

isApplyAgg

Boolean

是否应用聚合,为 false 时跳过该行

输出说明

返回类型

说明

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

事件计算公式,例如 event1/event2 表示 event1 除以 event2

sep

STRING

多值类别的分隔符,默认 \x1d

top_k

BIGINT

保留的最大键数,默认 50

is_punish

BOOLEAN

计算公式时是否用均值惩罚分母,默认 FALSE

输出说明

输出

类型

说明

kv_array

ARRAY

键值对数组,每个元素格式为 key1:value1\x1dkey2:value2,按事件类型和类别顺序排列

使用示例

-- 注册函数
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

pid

String

用户 ID

2

event_time

Long (bigint)

事件时间戳

3

fields

Array

需要记录的序列特征字段值数组

4

field_names

Array

序列特征字段名称数组,如 ["item_id", "category"]

5

event

String

当前事件类型

6

valid_events

Array

有效事件类型列表,如 ["click", "add", "buy"]

7

seq_len

Array

每个事件对应的序列长度限制

8

duplicate_method

Int

去重方法: - 1: 按 item+event 去重 - 2: 按 item+event+time 去重

9

combine_method

String

序列合并方法: - "all": 合并所有事件为一个序列 - 其他: 分别生成每个事件的序列

10

pre_seconds

Int

防止数据穿越的秒数,如 300 表示 5 分钟前的数据才生效

11

sequence_delim

String

序列分隔符,如 ";"

12

requstid

String

请求 ID,用于去重输出

13

scene

String

当前场景值

14

valid_scenes

Array

有效场景列表

参数说明:

  • args.length == 11: 无 requestid,无 scene

  • args.length == 12: 有 requestid,无 scene

  • args.length == 13: 无 requestid,有 scene

  • args.length == 14: 有 requestid,有 scene

输出说明

返回类型

说明

Map<String, String>

包含以下键值对: - "pid": 用户 ID - "event_time": 事件时间戳 - "request_id": 请求 ID (如果提供) - 序列特征: 格式为 {prefix}__{field_name},如 all_10_seq__item_id - 时间戳序列: 格式为 {prefix}__ts

输出序列键名格式

  • 如果 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 功能基本相同,但是去重逻辑略有差异。RtSeqFeature2duplicate_method=2 时,去重键为 fields[0] + event_time,而 RtSeqFeatureevent + fields[0] + event_time

入参说明

RtSeqFeature 完全相同。

输出说明

RtSeqFeature 完全相同。

与 RtSeqFeature 的区别

特性

RtSeqFeature

RtSeqFeature2

去重键 (duplicate_method=2)

event + fields[0] + event_time

fields[0] + event_time

去重键 (duplicate_method=1)

event + fields[0]

fields[0]

使用场景:

  • 如果需要根据 item 去重而不考虑事件类型,使用 RtSeqFeature2

  • 如果需要同时考虑事件类型和 item,使用 RtSeqFeature

3. MessageDelay

函数类com.aliyun.pai.udf.CountKV.MessageDelay
功能简介:根据事件时间和滑动窗口配置,计算消息的延迟时间。用于实时特征生成时,为不同时间窗口的统计数据生成延迟行。

入参说明

参数序号

参数名

类型

必填

说明

1

fields

Array

需要输出的基础字段数组

2

event_time

Long (bigint)

事件时间戳

3

windows

Array

时间窗口列表 (秒),如 [1800, 3600, 7200]

4

is_save

Boolean

是否保存 (生成延迟行): - true: 为每个窗口生成一行 - false: 只输出当前行

输出说明

返回类型

说明

Array

输出数组,包含: - 原始 fields 的所有字段 - 倒数第二个: 窗口触发时间戳 (event_time + window) - 最后一个:窗口大小 (秒)

使用示例

-- 为不同时间窗口生成延迟数据行
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

pid

String

用户 ID

2

cates

Array

类别/标签字段数组,如 ["style_id", "goods_id"]

3

cate_names

Array

类别字段名称数组,如 ["style_id", "goods_id"]

4

value

Int

统计值 (通常为 1)

5

event

String

当前事件类型

6

valid_events

Array

有效事件类型列表

7

delay_flag

Int

延迟标志: - 0: 立即生效 (添加到所有窗口) - 其他: 对应特定窗口 (用于减去过期数据)

8

windows

Array

时间窗口列表 (秒),如 [1800, 3600, 7200, 86400]

9

topk

Int

每个类别保留的 Top-K 数量

10

pre_seconds

Int

防止数据穿越的秒数

11

event_time

Long (bigint)

事件时间戳

12

output_date

String

输出日期 (格式: yyyyMMdd),只有大于等于该日期的数据才输出

13

side

String

字段归属标识,如 "user""item"

14

requstid

String

请求 ID,用于去重输出

15

scene

String

当前场景值

16

valid_scenes

Array

有效场景列表

参数说明:

  • args.length == 13:无 requestid,无 scene。

  • args.length == 14:有 requestid,无 scene。

  • args.length == 15:无 requestid,有 scene。

  • args.length == 16:有 requestid,有 scene。

输出说明

返回类型

说明

Map<String, String>

包含以下键值对: - "pid": 用户 ID - "event_time": 事件时间戳 - "request_id": 请求 ID (如果提供) - 事件计数: {side}__cnt_{event}_{window_name} - 类别键值对: {side}__kv_{cate_name}_{event}_{window_name}

输出键名格式

  • 窗口名称格式:rt{X}h{Y}m{Z}s,如 rt30mrt3hrt24h

  • 事件计数: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)的影响。

  • 适合推荐系统中稀有物品的质量评估。