PySpark如何基于日期列向DataFrame添加行
PySpark实现每行扩展为日期递增7天的多行记录
原始数据
首先定义初始的PySpark DataFrame:
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, array, lit, date_add, expr spark = SparkSession.builder.appName("DateIncrementDemo").getOrCreate() df = spark.createDataFrame( [(0, 0, 0, '2023-03-03'), (1, 1, 1, '2023-03-04'), (2, 2, 2, '2023-03-05')], ['item_id', 'action_id', 'status', 'date'] ) df.show()
原始输出:
+-------+---------+------+----------+ |item_id|action_id|status| date| +-------+---------+------+----------+ | 0| 0| 0|2023-03-03| | 1| 1| 1|2023-03-04| | 2| 2| 2|2023-03-05| +-------+---------+------+----------+
需求
为每一行生成3条重复行(加上原行共4条),每条的日期依次递增7天,得到目标结果。
解决方案
直接通过指定日期增量数组、展开数组并计算新日期的方式实现,步骤简洁高效:
# 1. 生成包含0、7、14、21天的增量数组,展开后每个增量对应一行 # 2. 将字符串日期转为Date类型,计算递增后的日期 # 3. 移除辅助列并排序 df_result = df.withColumn("days_increment", explode(array(lit(0), lit(7), lit(14), lit(21)))) \ .withColumn("date", date_add(expr("to_date(date)"), col("days_increment"))) \ .drop("days_increment") \ .orderBy("item_id", "date") df_result.show()
预期输出
+-------+---------+------+----------+ |item_id|action_id|status| date| +-------+---------+------+----------+ | 0| 0| 0|2023-03-03| | 0| 0| 0|2023-03-10| | 0| 0| 0|2023-03-17| | 0| 0| 0|2023-03-24| | 1| 1| 1|2023-03-04| | 1| 1| 1|2023-03-11| | 1| 1| 1|2023-03-18| | 1| 1| 1|2023-03-25| | 2| 2| 2|2023-03-05| | 2| 2| 2|2023-03-12| | 2| 2| 2|2023-03-19| | 2| 2| 2|2023-03-26| +-------+---------+------+----------+
说明
- 用
array(lit(0), lit(7), lit(14), lit(21))直接定义需要的日期增量,正好对应原日期+3次7天递增的需求。 explode函数将数组中的每个元素拆分为单独的行,实现每行扩展为4条记录。date_add配合to_date完成日期类型转换和递增计算,确保日期处理的正确性。
内容的提问来源于stack exchange,提问作者silkwire
相关产品推荐
相关产品推荐

