PySpark筛选满足多条件购买客户的实现问题求助
PySpark三类客户筛选完整解决方案
需求说明
需筛选满足以下任一条件的客户,并输出其customer_id和item列:
- 条件1:同时购买过
list_1和list_2中的商品 - 条件2:在
list_1范围内,同时购买过S/M和L/XL尺码的商品 - 条件3:在
list_2范围内,同时购买过S/M和L/XL尺码的商品
给定数据
商品列表
list_1 = ["A1", "A2", "B1", "B2", "C1", "C2", "D1", "D2"] list_2 = ["E1", "E2", "F1", "F2", "G1", "G2", "H1", "H2"]
尺码映射字典
dict_1 = {"A1" : "S/M", "A2" : "L/XL", "B1" : "S/M", "B2" : "L/XL", "C1" : "S/M", "C2" : "L/XL","D1" : "S/M", "D2" : "L/XL"} dict_2 = {"E1" : "S/M", "E2" : "L/XL", "F1" : "S/M", "F2" : "L/XL", "G1" : "S/M", "G2" : "L/XL", "H1" : "S/M", "H2" : "L/XL"}
原始数据表结构与示例数据
| customer_id | item | size |
|---|---|---|
| 1 | A1 | S/M |
| 2 | G1 | S/M |
| 3 | D1 | S/M |
| 1 | E2 | L/XL |
| 2 | H2 | L/XL |
| 9 | D1 | S/M |
| 1 | G1 | S/M |
| 9 | H2 | L/XL |
| 2 | H2 | L/XL |
已实现代码(条件1)
from pyspark.sql import functions as F from pyspark.sql.window import Window w = Window.partitionBy('customer_id').rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) df = (df.select('*', F.col('item').isin(list_1).alias('list_1'), F.col('item').isin(list_2).alias('list_2')) .select('customer_id', 'item', F.max('list_1').over(w).alias('list_1'), F.max('list_2').over(w).alias('list_2')) .filter(F.col('list_1') & F.col('list_2')) .select('customer_id', 'item'))
完整三类条件筛选代码
基于已有代码扩展,整合三类条件的判断逻辑:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口 w = Window.partitionBy('customer_id').rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) # 定义商品列表与尺码映射 list_1 = ["A1", "A2", "B1", "B2", "C1", "C2", "D1", "D2"] list_2 = ["E1", "E2", "F1", "F2", "G1", "G2", "H1", "H2"] dict_1 = {"A1" : "S/M", "A2" : "L/XL", "B1" : "S/M", "B2" : "L/XL", "C1" : "S/M", "C2" : "L/XL","D1" : "S/M", "D2" : "L/XL"} dict_2 = {"E1" : "S/M", "E2" : "L/XL", "F1" : "S/M", "F2" : "L/XL", "G1" : "S/M", "G2" : "L/XL", "H1" : "S/M", "H2" : "L/XL"} # 将字典转换为Spark可识别的映射表达式 map_list1_size = F.create_map([F.lit(x) for x in sum(dict_1.items(), ())]) map_list2_size = F.create_map([F.lit(x) for x in sum(dict_2.items(), ())]) # 计算客户满足的条件并筛选 df_result = df.select( "*", # 标记商品所属列表 F.col("item").isin(list_1).alias("in_list1"), F.col("item").isin(list_2).alias("in_list2") ).withColumn( # 条件1:同时购买过list1和list2商品 "cond1", F.max("in_list1").over(w) & F.max("in_list2").over(w) ).withColumn( # 条件2:list1范围内同时有S/M和L/XL购买记录 "cond2", F.max(F.when(F.col("in_list1") & (F.col("size") == "S/M"), 1)).over(w) >= 1 & F.max(F.when(F.col("in_list1") & (F.col("size") == "L/XL"), 1)).over(w) >= 1 ).withColumn( # 条件3:list2范围内同时有S/M和L/XL购买记录 "cond3", F.max(F.when(F.col("in_list2") & (F.col("size") == "S/M"), 1)).over(w) >= 1 & F.max(F.when(F.col("in_list2") & (F.col("size") == "L/XL"), 1)).over(w) >= 1 ).filter( # 满足任一条件即可 F.col("cond1") | F.col("cond2") | F.col("cond3") ).select( "customer_id", "item" ).distinct() # 去重避免重复记录 # 查看结果 df_result.show()
逻辑说明
- 映射转换:用
create_map将字典转为Spark可解析的映射,用于关联商品与对应尺码(也可直接使用原始size字段,映射用于双重校验) - 窗口计算:按
customer_id分组,计算每个客户是否满足三类条件:- 条件1:判断客户是否同时拥有list1和list2的购买记录
- 条件2/3:分别判断客户在对应商品列表中,是否同时存在两种尺码的购买记录
- 筛选去重:保留满足任一条件的记录,去重后得到目标数据表
内容的提问来源于stack exchange,提问作者Duc Vu
相关产品推荐
相关产品推荐

