PySpark按日期范围分组:无需UDF/Pandas实现组首行保留
PySpark 动态分组并保留每组首行实现方案
示例DataFrame
sample_df = spark.createDataFrame([ ('2020-01-01', '2021-01-01', 1), ('2020-02-01', '2021-02-01', 1), ('2021-01-15', '2022-01-15', 2), ('2022-01-15', '2023-01-15', 2), ('2022-02-01', '2023-02-01', 3), ('2022-03-01', '2023-03-01', 3), ('2023-03-01', '2024-03-01', 4), ], ['item_date', 'max_window', 'expected_grouping_index'])
需求说明
按
item_date排序后,首行启动一个分组;后续行若item_date≤当前组首行的max_window,则归为同组;若不在当前组范围内,则启动新分组。最终需保留每组的首行,要求实现过程不使用UDF,也不转换为Pandas DataFrame。
实现代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 转换日期字符串为日期类型,确保比较逻辑正确 df = sample_df.withColumn("item_date", F.to_date("item_date")) \ .withColumn("max_window", F.to_date("max_window")) # 定义排序窗口 sort_window = Window.orderBy("item_date") # 计算分组标识:判断当前行是否需要新建分组,累加得到分组ID df_with_group = df.withColumn("prev_group_max", F.lag("max_window").over(sort_window)) \ .withColumn("is_new_group", F.when(F.col("item_date") > F.col("prev_group_max"), 1).otherwise(0)) \ # 首行无前置窗口值,默认启动新分组 .withColumn("is_new_group", F.coalesce(F.col("is_new_group"), F.lit(1))) \ .withColumn("group_id", F.sum("is_new_group").over(sort_window) - 1) # 按分组ID取每组首行,还原原始字段并排序 result_df = df_with_group.groupBy("group_id") \ .agg( F.first("item_date").alias("item_date"), F.first("max_window").alias("max_window"), F.first("expected_grouping_index").alias("expected_grouping_index") ) \ .orderBy("item_date") # 输出结果 result_df.show()
逻辑说明
- 日期转换:将字符串类型的日期转为PySpark日期类型,避免字符串比较带来的逻辑错误。
- 分组标识计算:
- 用
lag函数获取上一个分组的max_window值; - 判断当前行的
item_date是否超出上一个分组的窗口范围,超出则标记为新分组; - 累加新分组标记得到唯一的
group_id,首行默认启动新分组。
- 用
- 保留首行:按
group_id分组后,用first函数提取每组的首行数据,最后按item_date排序还原顺序。
内容的提问来源于stack exchange,提问作者spatel4140
相关产品推荐
相关产品推荐

