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

基于客户操作顺序构建购物车物品数组的PySpark实现问题

问题描述

我做了个记录客户购物车操作的数据集,里面有ITEM_ADDED(加商品)、ITEM_REMOVED(删商品)、ITEM_UNDO(撤销操作)三类事件。举个例子:三次ITEM_ADDED之后连做三次ITEM_UNDO,购物车就空了。数据集如下:

event_idcustomer_ideventevent_typeevent_tsitem_idnext_event_type
246984993922{"item_id":1000,"customer_id":993922,"timestamp":5260,"type":"ITEM_ADDED"}ITEM_ADDED52601000ITEM_ADDED
246984993922{"item_id":1001,"customer_id":993922,"timestamp":5260,"type":"ITEM_ADDED"}ITEM_ADDED53551001ITEM_ADDED
246984993922{"item_id":1002,"customer_id":993922,"timestamp":5260,"type":"ITEM_ADDED"}ITEM_ADDED54601002ITEM_ADDED
246984993922{"after_id":0,"before_id":1002,"customer_id":993922,"timestamp":5500,"type":"ITEM_UNDO"}ITEM_UNDO5500NULLITEM_UNDO
246984993922{"after_id":0,"before_id":1001,"customer_id":993922,"timestamp":5510,"type":"ITEM_UNDO"}ITEM_UNDO5510NULLITEM_UNDO
246984993922{"after_id":0,"before_id":1000,"customer_id":993922,"timestamp":5515,"type":"ITEM_UNDO"}ITEM_UNDO5515NULLITEM_ADDED
246984993922{"item_id":2000,"customer_id":993922,"timestamp":6000,"type":"ITEM_ADDED"}ITEM_ADDED60002000ITEM_ADDED
246984993922{"item_id":2000,"customer_id":993922,"timestamp":6010,"type":"ITEM_REMOVED"}ITEM_ADDED60102000ITEM_REMOVED
246984993922{"item_id":7777,"customer_id":993922,"timestamp":9999,"type":"ITEM_ADDED"}ITEM_ADDED67007777ITEM_ADDED
246984993922{"item_id":9999,"customer_id":993922,"timestamp":9999,"type":"ITEM_ADDED"}ITEM_ADDED99999999NULL

预期最终购物车内容是[7777, 9999]。我试过用F.collect_list(F.col("item_id")).over(w.rowsBetween(Window.unboundedPreceding, end=0))来算,但这方法搞不定ITEM_REMOVED和ITEM_UNDO,求靠谱的实现方案。

我写的初始代码在这:

w = Window().partitionBy("customer_id").orderBy(F.asc("event_timestamp"))
cumulative_customer_items_purchased_df = (
    spark
    .table('testing')
    .where(
        (F.col("event_type").isin(item_events))
        &
        (F.col("customer_id") == 993922)
    )
    .select(*columns)
    .withColumn(
        "prev_event_type",
        F.lag(F.col("event_type")).over(w)
    )
    .withColumn(
        "next_event_type",
        F.lead(F.col("event_type")).over(w)
    )
)
解决方案

要搞定这种带撤销、移除操作的购物车状态计算,核心要么是跟踪每个商品的有效状态,要么用栈模拟操作顺序——毕竟ITEM_UNDO是撤销最近的操作,栈的逻辑最贴合业务。下面给两种可行方案:

方案1:用栈模拟操作(推荐,完全适配ITEM_UNDO逻辑)

直接用栈来模拟用户的每一步操作,逻辑清晰,不容易出错。我们可以写个UDF来维护每个用户的操作栈,一步步算出最终购物车。

from pyspark.sql import functions as F
from pyspark.sql.types import ArrayType, IntegerType

# 定义处理购物车操作的UDF
def process_cart_ops(ops):
    stack = []
    for op in ops:
        event_type = op["event_type"]
        item_id = op["item_id"]
        # 处理ITEM_UNDO:从event字段里取出要撤销的商品ID(before_id)
        if event_type == "ITEM_UNDO":
            before_id = op["event"]["before_id"]
            # 撤销最近的操作,也就是从栈里删掉这个商品
            if before_id in stack:
                stack.remove(before_id)
        elif event_type == "ITEM_ADDED":
            stack.append(item_id)
        elif event_type == "ITEM_REMOVED":
            # 直接移除指定商品
            if item_id in stack:
                stack.remove(item_id)
    return stack

# 注册UDF
process_cart_udf = F.udf(process_cart_ops, ArrayType(IntegerType()))

# 执行计算
final_cart_df = (
    spark.table('testing')
    .filter(F.col("customer_id") == 993922)
    # 按时间顺序把用户的所有操作打包成数组
    .groupBy("customer_id")
    .agg(F.collect_list(F.struct(
        "event_type", "item_id", "event"
    )).orderBy("event_ts").alias("operations"))
    # 用UDF处理操作数组,得到最终购物车
    .withColumn("final_cart", process_cart_udf(F.col("operations")))
    .select("customer_id", "final_cart")
)

# 查看结果
final_cart_df.show(truncate=False)

方案1说明

  • 先按用户分组,把所有操作按时间顺序打包成数组
  • 遍历操作数组,用栈模拟购物车变化:
    • 加商品就把ID塞栈里
    • 删商品就从栈里移除对应ID
    • 撤销操作就从event里拿到被撤销的商品ID,再从栈里删掉它
  • 最后栈里剩下的就是当前购物车的有效商品

方案2:用窗口函数标记有效操作

如果不想写UDF,也可以用窗口函数给每个操作打标记,筛选出没被抵消的添加操作。

from pyspark.sql import Window

w = Window.partitionBy("customer_id").orderBy("event_ts")

# 第一步:给每个操作编序号,解析ITEM_UNDO对应的目标商品
event_df = (
    spark.table('testing')
    .filter(F.col("customer_id") == 993922)
    .withColumn("op_id", F.row_number().over(w))
    # 从event字段里提取ITEM_UNDO要撤销的商品ID
    .withColumn("undo_target_item", F.when(F.col("event_type") == "ITEM_UNDO", 
                                           F.from_json(F.col("event"), "struct<before_id:int>").before_id)
               .otherwise(None))
)

# 第二步:标记每个添加操作是否被撤销或移除
# 反向窗口,检查当前商品之后有没有撤销/移除操作
w_reverse = Window.partitionBy("customer_id", "item_id").orderBy(F.desc("event_ts"))

marked_df = (
    event_df
    # 标记是否被ITEM_UNDO撤销
    .withColumn("is_undone", F.exists(
        F.collect_list(F.when(F.col("event_type") == "ITEM_UNDO", F.col("undo_target_item"))
                       .over(w_reverse)),
        lambda x: x == F.col("item_id")
    ))
    # 标记是否被ITEM_REMOVED移除
    .withColumn("is_removed", F.exists(
        F.collect_list(F.when(F.col("event_type") == "ITEM_REMOVED", F.col("item_id"))
                       .over(w_reverse)),
        lambda x: x == F.col("item_id")
    ))
    # 只保留没被撤销、没被移除的添加操作
    .filter((F.col("event_type") == "ITEM_ADDED") & ~F.col("is_undone") & ~F.col("is_removed"))
)

# 第三步:收集最终商品
final_cart_df = (
    marked_df.groupBy("customer_id")
    .agg(F.collect_list("item_id").alias("final_cart"))
)

final_cart_df.show(truncate=False)

方案2说明

  • 先给每个操作编个顺序号,解析出ITEM_UNDO要处理的商品ID
  • 用反向窗口函数,检查每个添加操作之后有没有对应的撤销或移除操作
  • 最后筛选出所有有效添加操作,收集商品ID就是最终购物车

结果验证

两种方案跑出来,customer_id=993922的最终购物车都是[7777, 9999],和预期一致。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 15:54:36