基于客户操作顺序构建购物车物品数组的PySpark实现问题
问题描述
我做了个记录客户购物车操作的数据集,里面有ITEM_ADDED(加商品)、ITEM_REMOVED(删商品)、ITEM_UNDO(撤销操作)三类事件。举个例子:三次ITEM_ADDED之后连做三次ITEM_UNDO,购物车就空了。数据集如下:
| event_id | customer_id | event | event_type | event_ts | item_id | next_event_type |
|---|---|---|---|---|---|---|
| 246984 | 993922 | {"item_id":1000,"customer_id":993922,"timestamp":5260,"type":"ITEM_ADDED"} | ITEM_ADDED | 5260 | 1000 | ITEM_ADDED |
| 246984 | 993922 | {"item_id":1001,"customer_id":993922,"timestamp":5260,"type":"ITEM_ADDED"} | ITEM_ADDED | 5355 | 1001 | ITEM_ADDED |
| 246984 | 993922 | {"item_id":1002,"customer_id":993922,"timestamp":5260,"type":"ITEM_ADDED"} | ITEM_ADDED | 5460 | 1002 | ITEM_ADDED |
| 246984 | 993922 | {"after_id":0,"before_id":1002,"customer_id":993922,"timestamp":5500,"type":"ITEM_UNDO"} | ITEM_UNDO | 5500 | NULL | ITEM_UNDO |
| 246984 | 993922 | {"after_id":0,"before_id":1001,"customer_id":993922,"timestamp":5510,"type":"ITEM_UNDO"} | ITEM_UNDO | 5510 | NULL | ITEM_UNDO |
| 246984 | 993922 | {"after_id":0,"before_id":1000,"customer_id":993922,"timestamp":5515,"type":"ITEM_UNDO"} | ITEM_UNDO | 5515 | NULL | ITEM_ADDED |
| 246984 | 993922 | {"item_id":2000,"customer_id":993922,"timestamp":6000,"type":"ITEM_ADDED"} | ITEM_ADDED | 6000 | 2000 | ITEM_ADDED |
| 246984 | 993922 | {"item_id":2000,"customer_id":993922,"timestamp":6010,"type":"ITEM_REMOVED"} | ITEM_ADDED | 6010 | 2000 | ITEM_REMOVED |
| 246984 | 993922 | {"item_id":7777,"customer_id":993922,"timestamp":9999,"type":"ITEM_ADDED"} | ITEM_ADDED | 6700 | 7777 | ITEM_ADDED |
| 246984 | 993922 | {"item_id":9999,"customer_id":993922,"timestamp":9999,"type":"ITEM_ADDED"} | ITEM_ADDED | 9999 | 9999 | NULL |
预期最终购物车内容是[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
相关产品推荐
相关产品推荐

