PySpark合并两个DataFrame并按时间过滤统计订单数的实现方法
错误原因分析
报错核心原因是代码中join操作后单独调用df1.filter(),此时df1本身不含df2的latest_order、关联后的addr_id属性,导致Spark找不到对应字段抛出异常。此外逻辑存在嵌套调用错误:select方法内不能直接调用count()方法返回整数值,需要按分组维度聚合统计。
正确实现方案
方案1:基于已有df2的关联统计
可直接复用已生成的df2,关联后过滤再聚合即可,代码如下:
import pyspark.sql.functions as f # 关联df1和df2,拿到每个订单对应的同用户同地址的最近下单时间 joined_df = df1.join(df2, on=["cust_id", "addr_id"], how="inner") # 过滤符合时间范围的订单后按分组统计数量 result = joined_df.filter( (f.col("order_time") >= f.date_sub(f.col("latest_order"), 30)) & (f.col("order_time") <= f.date_sub(f.col("latest_order"), 1)) ).groupBy("cust_id", "addr_id").agg( f.count("order_time").alias("order_counts"), f.first("latest_order").alias("latest_order") ) result.show()
方案2:窗口函数实现(性能更优,无需单独生成df2)
如果不需要单独保留df2,可直接用窗口函数一次性完成计算,避免多余的shuffle操作,适合大数据量场景:
import pyspark.sql.functions as f from pyspark.sql.window import Window # 定义窗口:按用户+地址分组计算每组的最大下单时间 w = Window.partitionBy("cust_id", "addr_id") result = df1.withColumn("latest_order", f.max("order_time").over(w)) \ .filter( (f.col("order_time") >= f.date_sub(f.col("latest_order"), 30)) & (f.col("order_time") <= f.date_sub(f.col("latest_order"), 1)) ).groupBy("cust_id", "addr_id").agg( f.count("order_time").alias("order_counts"), f.first("latest_order").alias("latest_order") ) result.show()
示例数据结果说明
给出的示例中df1所有订单时间均为2021-01-27,对比四个分组的latest_order,所有订单都不在「最近下单时间往前30天」的范围内,统计得到的order_counts均为0,实际业务中有符合时间范围的订单会正常统计对应数量。
内容的提问来源于stack exchange,提问作者Jitesh Malipeddi
相关产品推荐
相关产品推荐

