PySpark:基于活动DataFrame数组值筛选匹配的客户数据
PySpark 批量匹配活动与客户数据方案
问题背景
现有两张PySpark DataFrame:
活动表 (campaign_df)
| name|segment_list|rung_list | +--------------------+------------+-----------+ | Campaign 1 | [1.0, 5.0]| [L2, L3]| | Campaign 1 | [1.1]| [L1]| | Campaign 2 | [1.2]| [L2]| | Campaign 2 | [1.1]| [L4, L5]| +--------------------+------------+-----------+
客户表 (customer_df)
+-----------+---------------+---------+ |customer_id| segment |rung | +-----------+---------------+---------+ | 124001823| 1.0| L2| | 166001989| 5.0| L2| | 768002266| 1.1| L1| +-----------+---------------+---------+
需求说明
根据活动表的segment_list和rung_list匹配客户表数据,得到符合条件的活动-客户关联结果(示例仅保留Campaign 1的匹配记录),同时避免使用collect()循环逐行处理或UDF。
解决方案
使用PySpark内置的exists表达式实现条件匹配,通过内连接关联两张表,无需循环或UDF:
代码实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化SparkSession spark = SparkSession.builder.appName("CampaignCustomerMatch").getOrCreate() # 创建活动表 campaign_data = [ ("Campaign 1", [1.0, 5.0], ["L2", "L3"]), ("Campaign 1", [1.1], ["L1"]), ("Campaign 2", [1.2], ["L2"]), ("Campaign 2", [1.1], ["L4", "L5"]) ] campaign_df = spark.createDataFrame(campaign_data, ["name", "segment_list", "rung_list"]) # 创建客户表 customer_data = [ (124001823, 1.0, "L2"), (166001989, 5.0, "L2"), (768002266, 1.1, "L1") ] customer_df = spark.createDataFrame(customer_data, ["customer_id", "segment", "rung"]) # 匹配关联并生成结果 result_df = campaign_df.join( customer_df, # 检查客户的segment在活动的segment_list中,且rung在rung_list中 F.expr("exists(segment_list, x -> x = segment) AND exists(rung_list, x -> x = rung)"), how="inner" ).select("name", "customer_id") \ .distinct() # 去重,避免同一客户被同一活动的多个条件组重复匹配 # 如果仅需要Campaign 1的结果,添加过滤 result_df.filter(F.col("name") == "Campaign 1").show()
输出结果
+-----------+-----------+ | name|customer_id| +-----------+-----------+ |Campaign 1| 124001823| |Campaign 1| 166001989| |Campaign 1| 768002266| +-----------+-----------+
方案说明
- 利用
exists表达式直接在Spark分布式引擎层面完成数组包含检查,无需将数据拉取到Driver端(避免collect()的性能问题) - 内连接仅保留满足匹配条件的记录,性能优于交叉连接后过滤
distinct()用于去除同一客户被同一活动的多个条件组匹配产生的重复记录
内容的提问来源于stack exchange,提问作者Fizi
相关产品推荐
相关产品推荐

