You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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实现

核心思路

  1. 将date转为datetime类型,方便日期计算
  2. 为每行生成上一年对应日期的前后一周连续日期序列
  3. 用字典映射历史数据的(id1,id2,date)->value关系,快速匹配目标日期的value值
  4. 将匹配结果转置为指定列名的新列,合并回原数据集

代码实现

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实现

核心思路

  1. 转换日期类型,生成每行的唯一标识用于后续合并
  2. 计算上一年中心日期及前后一周区间,用sequence生成连续日期序列并展开
  3. 左连接历史数据,填充缺失value为0
  4. 通过透视将日期对应的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 04:37:06