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

如何在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(包含两端日期)
  • 禁止使用交叉连接(cross join)
  • 最终结果仅保留有有效指标的客户(即至少一个范围的曝光/转化数值不为0)
  • 期望输出格式:
client imp_Range1 con_Range1 imp_Range2 con_Range2
123    4          1          1          0  
124    0          0          3          1

解决方案

步骤说明

  1. 日期格式转换:将字符串格式的日期转为PySpark可比较的日期类型,确保范围判断准确
  2. 数据关联:通过client字段关联两张表,避免交叉连接
  3. 条件聚合统计:用条件函数分别计算两个日期范围内的曝光、转化总和
  4. 过滤有效客户:剔除所有指标均为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 17:40:44