SparkSQL实现累计数组聚合:避免UDF的高性能方案问询
高效实现累计追加next_item_id到数组的方案
原始数据集
| event_id | user_id | event | event_type | event_ts | item_id | next_event_type | next_item_id |
|---|---|---|---|---|---|---|---|
| 246984 | 993922 | {"item_id":1000,"user_id":993922,"timestamp":5260,"type":"ITEM_PURCHASED"} | ITEM_PURCHASED | 5260 | 1000 | ITEM_PURCHASED | 1001 |
| 246984 | 993922 | {"item_id":1001,"user_id":993922,"timestamp":5260,"type":"ITEM_PURCHASED"} | ITEM_PURCHASED | 5855 | 1001 | ITEM_PURCHASED | 1002 |
预期结果
| event_id | user_id | event | event_type | event_ts | item_id | next_event_type | next_item_id | next_item_set |
|---|---|---|---|---|---|---|---|---|
| 246984 | 993922 | {"item_id":1000,"user_id":993922,"timestamp":5260,"type":"ITEM_PURCHASED"} | ITEM_PURCHASED | 5260 | 1000 | ITEM_PURCHASED | 1001 | [1000, 1001] |
| 246984 | 993922 | {"item_id":1001,"user_id":993922,"timestamp":5260,"type":"ITEM_PURCHASED"} | ITEM_PURCHASED | 5855 | 1001 | ITEM_PURCHASED | 1002 | [1000, 1001, 1002] |
现有代码问题
当前查询存在两个关键问题:
- 引用了尚未生成的列
next_item_set,SQL执行时会报错,因为lag(next_item_set)中的列是当前select语句要创建的,无法提前引用 - 未按
user_id分区进行窗口计算,会导致跨用户累计数据,结果不符合预期
正确实现方案
使用内置窗口函数实现,全程无需UDF,适配大数据集性能需求:
WITH processed_data AS ( SELECT event_id, user_id, event, event_type, event_ts, item_id, -- 保留原有逻辑生成next_event_type和next_item_id LEAD(event_type) OVER (PARTITION BY user_id ORDER BY event_ts) AS next_event_type, LEAD(item_id) OVER (PARTITION BY user_id ORDER BY event_ts) AS next_item_id, -- 获取当前用户序列的第一个商品ID,作为数组初始值 FIRST_VALUE(item_id) OVER (PARTITION BY user_id ORDER BY event_ts) AS first_item, -- 累计收集从第一行到当前行的所有next_item_id COLLECT_LIST(next_item_id) OVER ( PARTITION BY user_id ORDER BY event_ts ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS collected_next_items FROM tableA ) SELECT event_id, user_id, event, event_type, event_ts, item_id, next_event_type, next_item_id, -- 合并初始商品ID和累计的next_item数组 ARRAY_CONCAT(ARRAY(first_item), collected_next_items) AS next_item_set FROM processed_data;
方案说明
FIRST_VALUE(item_id):按用户分区、时间排序,获取每个用户的第一个商品ID,作为累计数组的起始值COLLECT_LIST(next_item_id):按用户分区、时间排序,累计收集从第一行到当前行的所有next_item_idARRAY_CONCAT:将初始商品ID数组和累计的next_item数组合并,得到预期的next_item_set- 全部使用SQL内置函数,避免UDF带来的性能损耗,适合大规模数据集处理
内容的提问来源于stack exchange,提问作者satoshi
相关产品推荐
相关产品推荐

