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

PySpark分组聚合结合datediff标记入库10天未出现商品方案

PySpark 实现入库10天后无观测商品标记

问题描述

我是PySpark初学者。
需要识别自入库首日起10天后未再被观测到的商品,在DataFrame中新增列,符合该条件的商品赋值为1,其余商品赋值为0。
最初的实现思路:

  • 按product_id对数据分组,计算每组seen_date的最大值
  • 计算分组内import_date与max(seen_date)的日期差
  • 根据每组date_diff的取值生成新的标记列

报错代码与错误信息

最初编写的日期差计算代码运行报错:

from pyspark.sql.window import Window
from pyspark.sql import functions as F

w = (Window()
    .partitionBy(df.product_id)
    .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing))

df.withColumn("date_diff", F.datediff(F.max(F.from_unixtime(F.col("import_date")).over(w)), F.from_unixtime(F.col("seen_date"))))

报错信息:

AnalysisException: It is not allowed to use a window function inside an aggregate function. Please use the inner window function in a sub-query.

后续计划通过UDF基于date_diff字段生成not_seen标记列:

not_seen = udf(lambda x: 0 if x >10 else 1, IntegerType())
df = df.withColumn('not_seen', not_seen("date_diff"))

样例数据

样例数据生成代码如下:

columns = ["product_id","import_date", "seen_date"]
data = [("123", "2014-05-06", "2014-05-07"),
        ("123", "2014-05-06", "2014-06-11"),
        ("125", "2015-01-02", "2015-01-03"),
        ("125", "2015-01-02", "2015-01-04"),
        ("128", "2015-08-06", "2015-08-25")]
dfFromData2 = spark.createDataFrame(data).toDF(*columns)
dfFromData2 = dfFromData2.withColumn("import_date",F.unix_timestamp(F.col("import_date"),'yyyy-MM-dd'))
dfFromData2 = dfFromData2.withColumn("seen_date",F.unix_timestamp(F.col("seen_date"),'yyyy-MM-dd'))

生成的DataFrame结构:

+----------+-----------+----------+
|product_id|import_date| seen_date|
+----------+-----------+----------+
|       123| 1399334400|1399420800|
|       123| 1399334400|1402444800|
|       125| 1420156800|1420243200|
|       125| 1420156800|1420329600|
|       128| 1438819200|1440460800|
+----------+-----------+----------+

解决方案

原代码报错原因

原代码报错核心原因是窗口函数调用位置错误:把.over(w)写在了from_unixtime列之后,又在外层套了F.max()聚合,导致出现“聚合函数内嵌套窗口函数”的错误。另外计划使用的Python UDF需要在Python和JVM之间序列化数据,性能远低于Spark原生函数,非必要不推荐使用。

最优实现(无UDF,纯原生函数)

直接通过窗口函数计算每个商品的最大观测日期,再计算日期差,最后用when/otherwise生成标记列即可,代码如下:

from pyspark.sql import Window
from pyspark.sql import functions as F

# 定义分区窗口,默认全分区范围为无边界前后,无需显式写rowsBetween
w = Window.partitionBy("product_id")

result = dfFromData2.withColumn(
    # 计算每个商品的最大观测日期,转成date类型方便计算差值
    "max_seen_date", F.from_unixtime(F.max("seen_date").over(w), 'yyyy-MM-dd').cast("date")
).withColumn(
    # 转换入库日期为date类型
    "import_date_conv", F.from_unixtime("import_date", 'yyyy-MM-dd').cast("date")
).withColumn(
    "date_diff", F.datediff("max_seen_date", "import_date_conv")
).withColumn(
    # 日期差<=10则标记为1,否则为0,完全替代UDF
    "not_seen", F.when(F.col("date_diff") <= 10, 1).otherwise(0)
).drop("max_seen_date", "import_date_conv") # 删除计算用临时列

运行结果

执行后输出结果符合预期:

+----------+-----------+----------+---------+--------+
|product_id|import_date| seen_date|date_diff|not_seen|
+----------+-----------+----------+---------+--------+
|       123| 1399334400|1399420800|       36|       0|
|       123| 1399334400|1402444800|       36|       0|
|       125| 1420156800|1420243200|        2|       1|
|       125| 1420156800|1420329600|        2|       1|
|       128| 1438819200|1440460800|       19|       0|
+----------+-----------+----------+---------+--------+

结果验证:

  • 商品123最后一次观测距离入库36天,超过10天,标记为0
  • 商品125最后一次观测距离入库2天,不足10天,标记为1
  • 商品128最后一次观测距离入库19天,超过10天,标记为0

内容的提问来源于stack exchange,提问作者user_5

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 18:33:26