PySpark筛选同时购买两个列表商品的客户问题求助
问题分析与解决方案
原代码错误原因
- 字段丢失:先执行
select("customer_id").distinct()后,DataFrame已无item字段,后续where子句引用item必然触发报错。 - 逻辑错误:
list_1与list_2无重叠元素,单个item不可能同时属于两个列表,原条件(F.col("item").isin(list_1)) & (F.col("item").isin(list_2))永远为假,无法筛选出目标客户。
正确实现方法
方法1:分组聚合标记筛选
通过标记商品所属列表,分组后检查客户是否同时满足两个购买条件:
from pyspark.sql import functions as F list_1 = ["A", "B", "C", "D"] list_2 = ["E", "F", "G", "H"] # 标记每条记录是否属于目标列表 df_marked = df.withColumn( "has_list1", F.when(F.col("item").isin(list_1), 1).otherwise(0) ).withColumn( "has_list2", F.when(F.col("item").isin(list_2), 1).otherwise(0) ) # 按客户分组,验证是否同时购买过两个列表的商品 qualified_customers = df_marked.groupBy("customer_id").agg( F.max("has_list1").alias("has_list1"), F.max("has_list2").alias("has_list2") ).where((F.col("has_list1") == 1) & (F.col("has_list2") == 1)) # 关联原始数据,获取符合条件的客户所有记录 full_result = df.join(qualified_customers, on="customer_id", how="inner")
方法2:集合交集验证
利用集合交集功能,检查客户购买的商品集合是否与两个列表都有重叠:
from pyspark.sql import functions as F list_1 = ["A", "B", "C", "D"] list_2 = ["E", "F", "G", "H"] # 收集每个客户的所有购买商品集合 customer_item_sets = df.groupBy("customer_id").agg( F.collect_set("item").alias("purchased_items") ) # 筛选同时与两个列表有交集的客户 qualified_customers = customer_item_sets.where( F.size(F.array_intersect(F.col("purchased_items"), F.array(*list_1))) > 0 & F.size(F.array_intersect(F.col("purchased_items"), F.array(*list_2))) > 0 ) # 关联原始数据得到完整记录 full_result = df.join(qualified_customers, on="customer_id", how="inner")
内容的提问来源于stack exchange,提问作者Duc Vu
相关产品推荐
相关产品推荐

