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

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()

逻辑说明

  1. 日期转换:将字符串类型的日期转为PySpark日期类型,避免字符串比较带来的逻辑错误。
  2. 分组标识计算:
    • 用lag函数获取上一个分组的max_window值;
    • 判断当前行的item_date是否超出上一个分组的窗口范围,超出则标记为新分组;
    • 累加新分组标记得到唯一的group_id,首行默认启动新分组。
  3. 保留首行:按group_id分组后,用first函数提取每组的首行数据,最后按item_date排序还原顺序。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 04:32:34