PySpark实操:如何从DataFrame生成7天分段的日期范围DataFrame
问题描述
我有一个仅含单条记录的PySpark DataFrame,结构与数据如下:
spark_session_tbl_df.printSchema() spark_session_tbl_df.show()
输出结果:
root |-- strm: string (nullable = true) |-- acad_career: string (nullable = true) |-- session_code: string (nullable = true) |-- sess_begin_dt: timestamp (nullable = true) |-- sess_end_dt: timestamp (nullable = true) |-- census_dt: timestamp (nullable = true) +----+-----------+------------+-------------------+-------------------+-------------------+ |strm|acad_career|session_code| sess_begin_dt| sess_end_dt| census_dt| +----+-----------+------------+-------------------+-------------------+-------------------+ |2228| UGRD| 1|2022-08-20 00:00:00|2022-12-03 00:00:00|2022-09-19 00:00:00| +----+-----------+------------+-------------------+-------------------+-------------------+
我需要生成一个新的DataFrame,每行代表一个7天的日期区间,格式示例如下:
+-------------------+-------------------+ | sess_begin_dt| sess_end_dt| +-------------------+-------------------+ |2022-08-20 |2022-08-27 | +-------------------+-------------------+ |2022-08-28 |2022-09-04 | +-------------------+-------------------+ |2022-09-05 |2022-09-12 | +-------------------+-------------------+ |2022-09-13 |2022-09-20 | +-------------------+-------------------+ |2022-09-21 |2022-09-28 | +-------------------+-------------------+ ..... +-------------------+-------------------+ |2022-11-26 |2022-12-03 | +-------------------+-------------------+
我尝试了以下代码,但不确定是否正确引用了原DataFrame,也不确定是否有更合适的实现方式:
from pyspark.sql.functions import sequence, to_date, explode, col date_range_df = spark.sql("SELECT sequence(to_date('sess_begin_dt'), to_date('sess_end_dt'), interval 7 day) as date").withColumn("date", explode(col("date"))) date_range_df.show()
解决方案
你的代码存在两个核心问题:
- 使用
spark.sql()时,直接用字符串sess_begin_dt会被当作字面量,而非引用原DataFrame的列; - 现有代码只生成了起始日期序列,未计算每个区间的结束日期,且最后一个区间需匹配原DataFrame的
sess_end_dt,不能直接加7天。
以下是两种可行的实现方式:
方式一:DataFrame API实现(推荐)
直接基于原DataFrame操作,无需临时视图,逻辑更直观:
from pyspark.sql import functions as F from pyspark.sql.types import DateType # 提取原DataFrame中的起始、结束日期并转为Date类型 base_dates = spark_session_tbl_df.select( F.to_date("sess_begin_dt").alias("start_date"), F.to_date("sess_end_dt").alias("end_date") ).first() start_date = base_dates["start_date"] end_date = base_dates["end_date"] # 生成7天间隔的起始日期序列 date_range_df = spark.createDataFrame( [(d,) for d in sequence(start_date, end_date, 7)], ["interval_start"] ).withColumn("interval_start", F.col("interval_start").cast(DateType())) # 计算区间结束日期:若下一个起始日期超出原结束日期,则用原结束日期,否则加6天(保证区间为7天) date_range_df = date_range_df.withColumn( "interval_end", F.when( F.date_add(F.col("interval_start"), 7) > end_date, end_date ).otherwise( F.date_add(F.col("interval_start"), 6) ) ) # 重命名列以匹配需求格式 date_range_df = date_range_df.withColumnRenamed("interval_start", "sess_begin_dt")\ .withColumnRenamed("interval_end", "sess_end_dt") date_range_df.show()
方式二:SQL语法实现(需注册临时视图)
若偏好SQL写法,需先将原DataFrame注册为临时视图,再引用列:
from pyspark.sql import functions as F # 将原DataFrame注册为临时视图 spark_session_tbl_df.createOrReplaceTempView("session_table") # 生成日期序列并计算区间 date_range_df = spark.sql(""" SELECT date AS sess_begin_dt, CASE WHEN date_add(date, 7) > (SELECT to_date(sess_end_dt) FROM session_table) THEN (SELECT to_date(sess_end_dt) FROM session_table) ELSE date_add(date, 6) END AS sess_end_dt FROM ( SELECT explode(sequence( (SELECT to_date(sess_begin_dt) FROM session_table), (SELECT to_date(sess_end_dt) FROM session_table), interval 7 day )) AS date ) """) date_range_df.show()
关键说明
- 用
date_add(interval_start, 6)而非date_add(interval_start,7),是因为区间为闭区间(例如2022-08-20到2022-08-27包含7天),最后一个区间直接匹配原结束日期,避免超出范围; - 两种方式都先将timestamp类型转为date类型,确保日期格式符合需求。
内容的提问来源于stack exchange,提问作者Aztec619
相关产品推荐
相关产品推荐

