如何在PySpark DataFrame中对每行执行复杂处理计算订单占比
PySpark 实现方案
核心思路
核心逻辑分为三个模块:日期格式处理适配滑动窗口、90天滑动窗口统计分子/分母计数、关联客户全地址实现行扩展并输出标记字段,以下是可直接运行的完整实现:
步骤1:数据预处理
首先处理日期格式,并提取每个客户的唯一地址列表用于后续行扩展:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 假设原始DataFrame名为df,先将字符串日期转为日期格式和时间戳 df = df.withColumn("Order_Date", F.to_date(F.col("Order_Date"), "dd-MMM-yy")) \ .withColumn("Order_Date_ts", F.unix_timestamp(F.col("Order_Date"))) # 提取每个客户的所有唯一地址 cust_addr_distinct = df.select("Customer_ID", "Address_ID").distinct()
步骤2:滑动窗口计算分子、分母
定义两个滑动窗口分别统计地址维度订单数(分子)和客户全维度订单数(分母):
# 窗口1:按客户+地址分区,统计当前订单日期往前90天内该地址的订单数(分子) win_cust_addr = Window.partitionBy("Customer_ID", "Address_ID") \ .orderBy("Order_Date_ts") \ .rangeBetween(-90*86400, 0) # 窗口2:按客户分区,统计当前订单日期往前90天内该客户所有地址的总订单数(分母) win_cust = Window.partitionBy("Customer_ID") \ .orderBy("Order_Date_ts") \ .rangeBetween(-90*86400, 0) # 计算每个原始订单的计数指标,保留原始地址用于后续标记 df_with_cnt = df.withColumn("addr_order_cnt", F.count("Order_ID").over(win_cust_addr)) \ .withColumn("total_order_cnt", F.count("Order_ID").over(win_cust)) \ .withColumnRenamed("Address_ID", "original_Address_ID")
步骤3:行扩展并输出最终结果
将每个订单行关联客户全地址实现行扩展,计算占比和地址标记:
# 关联客户全地址,1行扩展为客户地址数量行 df_expand = df_with_cnt.join(cust_addr_distinct, on="Customer_ID", how="left") # 计算最终输出字段 df_result = df_expand.withColumn("Order_Share", F.col("addr_order_cnt") / F.col("total_order_cnt")) \ .withColumn("Is_original_address", F.when(F.col("Address_ID") == F.col("original_Address_ID"), F.lit(True)) .otherwise(F.lit(False))) \ .select("Customer_ID", "Address_ID", "Order_ID", "Order_Share", "Is_original_address")
结果验证
以示例中的Order_ID=7为例,输出结果符合预期:
- Cust_1、Addr_1对应Order_Share=3/7≈0.4286,Is_original_address=False
- Cust_1、Addr_2对应Order_Share=2/7≈0.2857,Is_original_address=False
- Cust_1、Addr_3对应Order_Share=2/7≈0.2857,Is_original_address=True
内容的提问来源于stack exchange,提问作者Jitesh Malipeddi
相关产品推荐
相关产品推荐

