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

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_iditemsize
1A1S/M
2G1S/M
3D1S/M
1E2L/XL
2H2L/XL
9D1S/M
1G1S/M
9H2L/XL
2H2L/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()

逻辑说明

  1. 映射转换:用create_map将字典转为Spark可解析的映射,用于关联商品与对应尺码(也可直接使用原始size字段,映射用于双重校验)
  2. 窗口计算:按customer_id分组,计算每个客户是否满足三类条件:
    • 条件1:判断客户是否同时拥有list1和list2的购买记录
    • 条件2/3:分别判断客户在对应商品列表中,是否同时存在两种尺码的购买记录
  3. 筛选去重:保留满足任一条件的记录,去重后得到目标数据表

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 18:24:29