如何在PySpark中按日期范围关联DataFrame并统计曝光转化指标
问题:PySpark 按客户统计指定日期范围内的曝光与转化量
原始数据
df1(客户日期范围表)
client date dateM1 dateM3 dateP3 123 9-2021 8-2021 6-2021 12-2021 124 8-2022 7-2022 5-2022 11-2022 125 2-2022 1-2022 11-2021 5-2022
df2(客户曝光转化表)
client date imp con 123 5-2021 2 0 123 6-2021 4 1 123 10-2021 1 0 124 11-2022 1 1 124 10-2022 2 0 125 8-2022 1 1
需求
日期格式为字符串M-yyyy,需完成以下操作:
- 统计每个客户在两个日期范围内的总曝光量(
imp)与转化量(con):- Range1:
dateM3至dateM1(包含两端日期) - Range2:
date至dateP3(包含两端日期)
- Range1:
- 禁止使用交叉连接(cross join)
- 最终结果仅保留有有效指标的客户(即至少一个范围的曝光/转化数值不为0)
- 期望输出格式:
client imp_Range1 con_Range1 imp_Range2 con_Range2 123 4 1 1 0 124 0 0 3 1
解决方案
步骤说明
- 日期格式转换:将字符串格式的日期转为PySpark可比较的日期类型,确保范围判断准确
- 数据关联:通过
client字段关联两张表,避免交叉连接 - 条件聚合统计:用条件函数分别计算两个日期范围内的曝光、转化总和
- 过滤有效客户:剔除所有指标均为0的客户
代码实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化SparkSession spark = SparkSession.builder.appName("date_range_stats").getOrCreate() # 构造df1数据 data1 = [ (123, "9-2021", "8-2021", "6-2021", "12-2021"), (124, "8-2022", "7-2022", "5-2022", "11-2022"), (125, "2-2022", "1-2022", "11-2021", "5-2022") ] df1 = spark.createDataFrame(data1, ["client", "date", "dateM1", "dateM3", "dateP3"]) # 构造df2数据 data2 = [ (123, "5-2021", 2, 0), (123, "6-2021", 4, 1), (123, "10-2021", 1, 0), (124, "11-2022", 1, 1), (124, "10-2022", 2, 0), (125, "8-2022", 1, 1) ] df2 = spark.createDataFrame(data2, ["client", "date", "imp", "con"]) # 定义日期转换逻辑:将"M-yyyy"转为标准日期(取当月1号) def convert_date(date_col): return F.to_date(F.concat(F.lit("1-"), date_col), "d-M-yyyy") # 转换df1的所有日期列 df1 = df1.withColumn("date", convert_date("date")) \ .withColumn("dateM1", convert_date("dateM1")) \ .withColumn("dateM3", convert_date("dateM3")) \ .withColumn("dateP3", convert_date("dateP3")) # 转换df2的日期列并重命名,避免关联后字段冲突 df2 = df2.withColumn("event_date", convert_date("date")) # 按client关联两张表,保留df1所有客户 joined_df = df1.join(df2, on="client", how="left") # 按client聚合,统计两个范围的指标 result_df = joined_df.groupBy("client") \ .agg( # 统计Range1的曝光和转化 F.sum(F.when(F.col("event_date").between(F.col("dateM3"), F.col("dateM1")), F.col("imp")).otherwise(0)).alias("imp_Range1"), F.sum(F.when(F.col("event_date").between(F.col("dateM3"), F.col("dateM1")), F.col("con")).otherwise(0)).alias("con_Range1"), # 统计Range2的曝光和转化 F.sum(F.when(F.col("event_date").between(F.col("date"), F.col("dateP3")), F.col("imp")).otherwise(0)).alias("imp_Range2"), F.sum(F.when(F.col("event_date").between(F.col("date"), F.col("dateP3")), F.col("con")).otherwise(0)).alias("con_Range2") ) # 过滤掉所有指标都为0的客户 result_df = result_df.filter( (F.col("imp_Range1") != 0) | (F.col("con_Range1") != 0) | (F.col("imp_Range2") != 0) | (F.col("con_Range2") != 0) ) # 打印结果 result_df.show()
代码解释
- 日期转换:通过拼接"1-"将
M-yyyy转为d-M-yyyy格式,确保日期可以被PySpark正确解析和比较 - 关联逻辑:使用
left join按client关联,既保证每个客户的日期范围能匹配到对应的曝光数据,又避免了交叉连接的性能问题 - 条件聚合:用
when...otherwise判断事件日期是否在目标范围内,对符合条件的imp和con求和,不符合的记为0 - 过滤规则:通过逻辑或判断保留至少有一个有效指标的客户,满足需求中的结果要求
内容的提问来源于stack exchange,提问作者hbabbar
相关产品推荐
相关产品推荐

