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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 11:20:34