如何在PySpark/Pandas中按日期范围取数并转置关联上一年同期数据
问题描述
现有如下格式的数据集:
id1 id2 date value 0 1 2021-12-28 24 0 1 2021-12-30 24 0 1 2022-01-04 24 0 1 2022-01-06 24 0 1 2022-01-07 8 0 1 2022-01-11 16 0 1 2022-01-13 16 0 1 2023-01-03 16 0 1 2023-01-05 56
需求:针对每个id1与id2的组合,对每行的date字段,提取上一年中该日期前一周至后一周范围内的value列值,将这些值拆分为新列;若该范围内存在缺失日期,对应列需保留位置并填充0,且列需按顺序排列。
例如:当id1=0、id2=1、date=2023-01-03时,需提取id1=0、id2=1且date在2022-12-27至2022-01-09之间的value值并转置为新列,最终目标数据集格式如下:
id1 id2 date value day-7 day-6 day-5 day-4 day-3 day-2 day-1 day-0 day+1 day+2 day+3 day+4 day+5 day+6 day+7 0 1 2021-12-28 24 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 1 2021-12-30 24 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 1 2022-01-04 24 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 1 2022-01-06 24 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 1 2022-01-07 8 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 1 2022-01-11 16 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 1 2022-01-13 16 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 1 2023-01-03 16 0 24 0 24 0 0 0 0 24 0 24 8 0 0 0 0 1 2023-01-05 56 0 24 0 0 0 0 24 0 24 8 0 0 0 16 0
尝试过窗口函数和连接操作但未解决问题,求在PySpark或Pandas中的实现方案。
解决方案
一、Pandas实现
核心思路
- 将
date转为datetime类型,方便日期计算 - 为每行生成上一年对应日期的前后一周连续日期序列
- 用字典映射历史数据的
(id1,id2,date)->value关系,快速匹配目标日期的value值 - 将匹配结果转置为指定列名的新列,合并回原数据集
代码实现
import pandas as pd # 读取示例数据 df = pd.DataFrame({ 'id1': [0]*9, 'id2': [1]*9, 'date': ['2021-12-28', '2021-12-30', '2022-01-04', '2022-01-06', '2022-01-07', '2022-01-11', '2022-01-13', '2023-01-03', '2023-01-05'], 'value': [24,24,24,24,8,16,16,16,56] }) # 转换日期格式 df['date'] = pd.to_datetime(df['date']) # 构建历史数据映射表,提升查询效率 history_map = df.set_index(['id1', 'id2', 'date'])['value'].to_dict() # 定义每行处理函数 def extract_history(row): # 计算上一年的中心日期 target_center = row['date'] - pd.DateOffset(years=1) # 生成前后一周的连续日期 date_range = pd.date_range(target_center - pd.Timedelta(days=7), target_center + pd.Timedelta(days=7), freq='D') # 匹配每个日期的value,缺失则填0 values = [history_map.get((row['id1'], row['id2'], dt), 0) for dt in date_range] # 生成列名:day-7到day+7 cols = [f'day{"+" if i >=0 else ""}{i-7}' for i in range(15)] return pd.Series(values, index=cols) # 应用函数并合并结果 result = pd.concat([df, df.apply(extract_history, axis=1)], axis=1) # 调整列顺序匹配目标格式 target_cols = ['id1', 'id2', 'date', 'value'] + [f'day{"+" if i >=0 else ""}{i}' for i in range(-7, 8)] result = result[target_cols] print(result.to_string(index=False))
二、PySpark实现
核心思路
- 转换日期类型,生成每行的唯一标识用于后续合并
- 计算上一年中心日期及前后一周区间,用
sequence生成连续日期序列并展开 - 左连接历史数据,填充缺失value为0
- 通过透视将日期对应的value转置为指定列名的新列,合并回原数据集
代码实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import DateType spark = SparkSession.builder.appName("date_history_extract").getOrCreate() # 读取示例数据 data = [ (0, 1, "2021-12-28", 24), (0, 1, "2021-12-30", 24), (0, 1, "2022-01-04", 24), (0, 1, "2022-01-06", 24), (0, 1, "2022-01-07", 8), (0, 1, "2022-01-11", 16), (0, 1, "2022-01-13", 16), (0, 1, "2023-01-03", 16), (0, 1, "2023-01-05", 56) ] df = spark.createDataFrame(data, ["id1", "id2", "date", "value"]) df = df.withColumn("date", F.col("date").cast(DateType())) # 生成每行唯一标识,用于合并 df = df.withColumn("row_id", F.monotonically_increasing_id()) # 计算目标日期区间 df_with_range = df.withColumn( "target_center", F.add_months(F.col("date"), -12) ).withColumn( "start_date", F.date_sub(F.col("target_center"), 7) ).withColumn( "end_date", F.date_add(F.col("target_center"), 7) ) # 展开连续日期并计算偏移量 df_exploded = df_with_range.withColumn( "target_date", F.explode(F.sequence(F.col("start_date"), F.col("end_date"), F.expr("interval 1 day"))) ).withColumn( "day_offset", F.datediff(F.col("target_date"), F.col("target_center")) ).withColumn( "col_name", F.concat(F.lit("day"), F.when(F.col("day_offset") >=0, F.lit("+")), F.col("day_offset").cast("string")) ) # 左连接历史数据,填充缺失值 history_df = df.select("id1", "id2", "date", "value").withColumnRenamed("date", "target_date").withColumnRenamed("value", "history_value") df_joined = df_exploded.join(history_df, on=["id1", "id2", "target_date"], how="left").fillna(0, subset=["history_value"]) # 透视转宽表 pivot_df = df_joined.groupBy("row_id").pivot("col_name").agg(F.first("history_value")) # 合并并调整列顺序 result = df.join(pivot_df, on="row_id").drop("row_id", "target_center", "start_date", "end_date") target_cols = ["id1", "id2", "date", "value"] + [f"day{'+' if i >=0 else ''}{i}" for i in range(-7, 8)] result = result.select(target_cols) result.show(truncate=False)
内容的提问来源于stack exchange,提问作者abtExp
相关产品推荐
相关产品推荐

