PySpark中如何按起止日期生成对应月份的新行?
PySpark生成日期区间内逐月拆分数据并添加月份序号
核心思路
不用array_repeat机械重复行,而是通过生成月份偏移序列,结合add_months函数实现逐月日期递增,同时直接用偏移量+1得到从1开始的月份序号。
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import col, sequence, explode, add_months, lit, ceil, months_between from pyspark.sql.types import DateType # 初始化SparkSession spark = SparkSession.builder.appName("MonthSplitDemo").getOrCreate() # 1. 准备测试数据(替换为你的业务数据) data = [ ("user1", "2023-01-01", "2023-08-01"), ("user2", "2022-05-10", "2023-01-10") ] df = spark.createDataFrame(data, ["user_id", "start_date_str", "end_date_str"]) # 2. 转换字符串日期为Date类型 df = df.withColumn("start_date", col("start_date_str").cast(DateType())) \ .withColumn("end_date", col("end_date_str").cast(DateType())) # 3. 动态计算日期区间的总月份数(兼容任意间隔) # months_between返回小数,ceil取整后+1确保包含首尾月份 df = df.withColumn("total_months", ceil(months_between(col("end_date"), col("start_date"))).cast("int") + 1) # 4. 生成0到total_months-1的偏移序列,用于逐月递增 df = df.withColumn("month_offsets", sequence(lit(0), col("total_months") - 1)) # 5. 炸开序列,得到每行对应一个月份偏移量 df_exploded = df.withColumn("offset", explode(col("month_offsets"))) # 6. 生成目标月份日期和月份序号 result_df = df_exploded.withColumn("target_month", add_months(col("start_date"), col("offset"))) \ .withColumn("month_order", col("offset") + 1) \ .select("user_id", "start_date", "end_date", "target_month", "month_order") # 查看结果 result_df.orderBy("user_id", "month_order").show()
关键说明
- 为什么
array_repeat不符合预期?它只是机械重复指定次数的行,无法生成逐月递增的日期,也没法关联对应序号。 sequence生成的偏移量从0开始,配合add_months能精准得到每个月份的日期,offset + 1直接对应从1开始的月份序号。- 动态计算
total_months时,用ceil(months_between(...)) + 1是为了确保包含起始和结束月份(比如2023-01-01到2023-08-01,months_between返回7.0,ceil后+1得到8,对应8行数据)。
内容的提问来源于stack exchange,提问作者Tuhin
相关产品推荐
相关产品推荐

