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

如何按日期范围将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)

关键细节说明

  1. 周周期自定义:通过修改WEEK_END_DAY可以调整周的结束日,比如设置为"sun"时,周周期为周日到周六,可根据业务需求灵活调整周起始逻辑。
  2. 首尾段对齐:
    • 首段使用least函数确保不会超出原end_date,如果原start_date本身就是完整周起始,则首段自动转为完整周。
    • 末段使用greatest函数确保不会早于原start_date,如果原end_date本身就是完整周结束,则末段自动转为完整周。
  3. 无效分段过滤:当原日期范围本身就是一个完整周或更短时,过滤掉重复或逻辑无效的行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 18:29:59