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
相关产品推荐
相关产品推荐

