如何按日期范围将PySpark DataFrame行拆分为每周多行?
PySpark实现日期范围按周拆分(首尾对齐原日期)
核心思路
通过生成日期序列、构造周分段结构并拆分多行,确保首段起始对齐原start_date、末段结束对齐原end_date,中间为完整周周期。
完整代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import ( col, explode, sequence, date_add, next_day, least, greatest, lit, array, struct ) from pyspark.sql.types import DateType # 初始化SparkSession spark = SparkSession.builder.appName("WeekSplit").getOrCreate() # 创建测试DataFrame data = [ ("0001", "1111", "2019-10-10", "2019-11-06"), ("0002", "1112", "2020-11-26", "2021-01-06"), ("0003", "1113", "2020-09-24", "2020-10-21") ] df = spark.createDataFrame(data, ["client", "family_code", "start_date", "end_date"]) # 转换字符串日期为Date类型 df = df.withColumn("start_date", col("start_date").cast(DateType())) \ .withColumn("end_date", col("end_date").cast(DateType())) # 定义周结束日(可根据需求调整,比如"sun"表示周日为周结束) WEEK_END_DAY = "sat" # 生成周分段并拆分多行 df_with_weeks = df.withColumn( # 计算第一个完整周的起始日(示例:周日为周起始) "first_full_week_start", next_day(col("start_date"), "sun") ).withColumn( # 计算最后一个完整周的结束日 "last_full_week_end", date_add(next_day(col("end_date"), WEEK_END_DAY), -7) ).withColumn( # 生成中间完整周的起始日期序列 "middle_week_starts", sequence( col("first_full_week_start"), col("last_full_week_end"), lit(7) ) ).withColumn( # 构建所有周分段的结构列表 "week_segments", array( # 首段:从原start_date到第一个完整周的前一天 struct( col("start_date").alias("week_start"), least(date_add(col("first_full_week_start"), -1), col("end_date")).alias("week_end") ), # 中间完整周:每个起始日对应7天周期 explode( col("middle_week_starts").transform( lambda s: struct( s.alias("week_start"), date_add(s, 6).alias("week_end") ) ) ), # 末段:从最后一个完整周次日到原end_date struct( greatest(date_add(col("last_full_week_end"), 1), col("start_date")).alias("week_start"), col("end_date").alias("week_end") ) ) ).withColumn( # 拆分为多行 "week_segment", explode(col("week_segments")) ).select( "client", "family_code", col("week_segment.week_start").alias("week_start_date"), col("week_segment.week_end").alias("week_end_date"), "start_date", "end_date" ).filter( # 过滤无效分段(避免日期范围过短导致的重复/无效行) col("week_start_date") <= col("week_end_date") ) # 查看结果 df_with_weeks.orderBy("client", "week_start_date").show(truncate=False)
关键细节说明
- 周周期自定义:通过修改
WEEK_END_DAY可以调整周的结束日,比如设置为"sun"时,周周期为周日到周六,可根据业务需求灵活调整周起始逻辑。 - 首尾段对齐:
- 首段使用
least函数确保不会超出原end_date,如果原start_date本身就是完整周起始,则首段自动转为完整周。 - 末段使用
greatest函数确保不会早于原start_date,如果原end_date本身就是完整周结束,则末段自动转为完整周。
- 首段使用
- 无效分段过滤:当原日期范围本身就是一个完整周或更短时,过滤掉重复或逻辑无效的行。
内容的提问来源于stack exchange,提问作者jabeono
相关产品推荐
相关产品推荐

