基于起止日期列生成周结束日期列的PySpark问题
PySpark生成起止日期间的每周结束日期
原始数据集
ID start_date end_date 1 1999-04-07 2013-09-15 1 2015-08-15 2017-02-01 2 2018-05-04 2020-04-02 8 2000-01-03 2011-07-19
预期效果
ID start_date end_date Weekly_end_date 1 1999-04-07 2013-09-15 1999-04-13 1 1999-04-07 2013-09-15 1999-04-20 1 1999-04-07 2013-09-15 1999-04-27 ... 1 1999-04-07 2013-09-15 2013-09-15 1 2015-08-15 2017-02-01 2015-08-22 ... 1 2015-08-15 2017-02-01 2017-02-01 2 2018-05-04 2020-04-02 2018-05-05 ... 2 2018-05-04 2020-04-02 2020-04-02 8 2000-01-03 2011-07-19 2000-01-08 ... 8 2000-01-03 2011-07-19 2011-07-19
问题分析
你当前的代码直接用sequence('start_date', 'end_date', interval 1 week)生成的是从start_date当天开始每周的同一天,而非每周结束日期;同时若start_date到end_date间隔不足一周,会出现无数据生成的情况,也无法保证end_date被纳入结果。
解决方案
步骤说明
- 计算首个每周结束日期:按示例以周六作为周结束(可按需调整),通过
dayofweek判断当前日期是周几,计算得到下一个周六;若当天已是周六则直接保留。 - 生成完整日期序列:从首个周结束日期开始,按每周间隔生成序列,同时判断
end_date是否在序列中,若不在则追加该日期,确保结果包含终止日期。 - 展开序列:用
explode将数组列展开为多行,得到每行对应的每周结束日期。
完整代码
from pyspark.sql import functions as F from pyspark.sql.types import DateType # 创建原始DataFrame(实际场景可替换为读取数据源) df = spark.createDataFrame( [ (1, "1999-04-07", "2013-09-15"), (1, "2015-08-15", "2017-02-01"), (2, "2018-05-04", "2020-04-02"), (8, "2000-01-03", "2011-07-19") ], ["ID", "start_date", "end_date"] ).withColumn("start_date", F.col("start_date").cast(DateType()))\ .withColumn("end_date", F.col("end_date").cast(DateType())) # 计算第一个每周结束日期(周六为周结束) df = df.withColumn( "first_weekly_end", F.when( F.dayofweek(F.col("start_date")) == 7, # 当天是周六则直接保留 F.col("start_date") ).otherwise( F.date_add(F.col("start_date"), 7 - F.dayofweek(F.col("start_date"))) ) ) # 生成包含end_date的每周结束日期序列 df = df.withColumn( "weekly_end_dates", F.when( F.col("first_weekly_end") > F.col("end_date"), F.array(F.col("end_date")) # 首个结束日期晚于终止日,直接用终止日 ).otherwise( F.array_union( F.sequence(F.col("first_weekly_end"), F.col("end_date"), F.expr("interval 1 week")), F.when( F.element_at(F.sequence(F.col("first_weekly_end"), F.col("end_date"), F.expr("interval 1 week")), -1) != F.col("end_date"), F.array(F.col("end_date")) ).otherwise(F.array()) ) ) ) # 展开序列并清理临时列 result_df = df.withColumn("weekly_end_date", F.explode(F.col("weekly_end_dates")))\ .drop("first_weekly_end", "weekly_end_dates") # 查看结果 result_df.show()
自定义调整说明
- 周结束日期切换:若需将周日作为周结束,只需修改判断逻辑:将
dayofweek(start_date) == 7改为dayofweek(start_date) == 1,计算首个结束日期的公式改为date_add(start_date, 1 - dayofweek(start_date))。 - 边界场景处理:当
start_date的下一个周结束日期晚于end_date时,直接将end_date作为唯一的周结束日期输出。
内容的提问来源于stack exchange,提问作者MS128
相关产品推荐
相关产品推荐

