如何在PySpark中按Page和Group将DataFrame转为带2天滞后的时间序列
PySpark生成时间序列滞后特征的实现方法
要实现按page和group分组,生成t的1天和2天滞后特征,核心是用PySpark的窗口函数结合lag函数,具体操作如下:
步骤1:导入必要函数
首先导入窗口函数和lag函数:
from pyspark.sql import Window from pyspark.sql.functions import lag
步骤2:定义窗口规范
创建窗口,指定分组键(page和group)以及排序规则(按utc_date升序),确保lag函数在每个分组内按时间顺序计算滞后值:
window_spec = Window.partitionBy("page", "group").orderBy("utc_date")
步骤3:生成滞后特征列
使用withColumn方法添加t-1和t-2列,分别对应滞后1期和2期的t值:
result_df = df.withColumn("t-1", lag("t", 1).over(window_spec)) \ .withColumn("t-2", lag("t", 2).over(window_spec))
步骤4:查看结果
执行show()方法即可得到目标格式的数据集:
result_df.show()
关键说明
partitionBy("page", "group"):将计算范围限定在每个page+group的组合内,不同分组之间不会互相干扰。orderBy("utc_date"):保证时间序列按日期顺序排列,让lag函数能正确取到前N天的数值。lag("t", n):参数n代表滞后的行数,由于数据按天连续,滞后1行就是前1天的t值;没有前置行时会返回null,完全匹配需求格式。
内容的提问来源于stack exchange,提问作者Munichong
相关产品推荐
相关产品推荐

