PySpark如何用内置函数获取用户近5天内最后3件购买商品?
使用PySpark内置函数提取用户近5天内最后3件商品
完全可以用PySpark内置函数实现,无需自定义UDF,核心思路是通过窗口函数筛选目标数据,再聚合收集结果,具体实现如下:
实现步骤
- 数据格式转换:将字符串类型的
timestamp转为日期类型,方便后续日期范围筛选。 - 筛选近5天记录:基于当前日期(或数据集内的最大日期),过滤出用户近5天的购买数据。
- 窗口排序编号:按用户
id分区,按timestamp降序(最新记录优先)给每条记录编号;若同一天有多条记录,用自增ID保留原始顺序。 - 提取前3条数据:筛选每个用户分区内编号≤3的记录。
- 聚合生成结果:按
id分组,将筛选后的item收集为列表,保持排序后的顺序。
代码示例
假设原始数据集已加载为DataFramedf,代码如下:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("LastThreeItems").getOrCreate() # 构建示例数据 data = [ (1, "apple", "2022-08-15"), (1, "peach", "2022-08-15"), (1, "apple", "2022-08-15"), (1, "banana", "2022-08-14"), (2, "apple", "2022-08-15"), (2, "banana", "2022-08-14"), (2, "banana", "2022-08-14"), (2, "water", "2022-08-14"), (3, "water", "2022-08-15"), (3, "water", "2022-08-14") ] df = spark.createDataFrame(data, ["id", "item", "timestamp"]) # 转换日期格式并筛选近5天数据 df = df.withColumn("date", F.to_date(F.col("timestamp"))) # 若以数据集内最新日期为基准,可替换为:max_date = df.select(F.max("date")).first()[0] filtered_df = df.filter(F.col("date") >= F.date_sub(F.current_date(), 5)) # 定义窗口规则 window_spec = Window.partitionBy("id").orderBy(F.col("date").desc(), F.monotonically_increasing_id()) # 生成结果 result_df = filtered_df.withColumn("row_num", F.row_number().over(window_spec)) \ .filter(F.col("row_num") <= 3) \ .groupBy("id") \ # 若需将最新商品放在列表末尾,替换为F.reverse(F.collect_list("item")) .agg(F.collect_list("item").alias("last_three_item")) \ .orderBy("id") # 展示结果 result_df.show(truncate=False)
补充说明
- 如果需要以数据集内的最新日期作为近5天的计算基准,只需替换筛选条件中的
F.current_date()为数据集的最大日期即可。 - 若对列表内的商品顺序有特殊要求(如用户示例中id2的倒序展示),可在聚合时使用
F.reverse()函数反转收集到的列表。
内容的提问来源于stack exchange,提问作者DutchBro
相关产品推荐
相关产品推荐

