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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 17:36:05