You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

SparkSQL实现累计数组聚合:避免UDF的高性能方案问询

高效实现累计追加next_item_id到数组的方案

原始数据集

event_iduser_ideventevent_typeevent_tsitem_idnext_event_typenext_item_id
246984993922{"item_id":1000,"user_id":993922,"timestamp":5260,"type":"ITEM_PURCHASED"}ITEM_PURCHASED52601000ITEM_PURCHASED1001
246984993922{"item_id":1001,"user_id":993922,"timestamp":5260,"type":"ITEM_PURCHASED"}ITEM_PURCHASED58551001ITEM_PURCHASED1002

预期结果

event_iduser_ideventevent_typeevent_tsitem_idnext_event_typenext_item_idnext_item_set
246984993922{"item_id":1000,"user_id":993922,"timestamp":5260,"type":"ITEM_PURCHASED"}ITEM_PURCHASED52601000ITEM_PURCHASED1001[1000, 1001]
246984993922{"item_id":1001,"user_id":993922,"timestamp":5260,"type":"ITEM_PURCHASED"}ITEM_PURCHASED58551001ITEM_PURCHASED1002[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_id
  • ARRAY_CONCAT:将初始商品ID数组和累计的next_item数组合并,得到预期的next_item_set
  • 全部使用SQL内置函数,避免UDF带来的性能损耗,适合大规模数据集处理

内容的提问来源于stack exchange,提问作者satoshi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.23 15:57:18