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

在PySpark DataFrame中新增列:计算指定日期区间价格最大值

PySpark实现日期区间内价格最大值计算

现有DataFrame

dateprice
1995-10-2410
1995-10-2520
1995-10-2630
1995-10-2711
1995-10-2814
1995-10-2912

需求说明

新增一列new_price,取值为当前日期与该日期前2天区间内的price最大值,即new_price = max(price) between (date, date-2 days)。预期结果如下:

预期结果

datepricenew_price
1995-10-241010
1995-10-252020
1995-10-263030
1995-10-271130
1995-10-281430
1995-10-291214

代码实现

方法一:基于时间戳的区间范围(推荐,适配日期不连续场景)

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import col, max, unix_timestamp

# 初始化SparkSession
spark = SparkSession.builder.appName("DateRangeMaxPrice").getOrCreate()

# 创建示例数据
data = [
    ("1995-10-24", 10),
    ("1995-10-25", 20),
    ("1995-10-26", 30),
    ("1995-10-27", 11),
    ("1995-10-28", 14),
    ("1995-10-29", 12)
]
df = spark.createDataFrame(data, ["date", "price"])

# 将字符串日期转为日期类型,并生成unix时间戳列
df = df.withColumn("date", col("date").cast("date")) \
       .withColumn("unix_date", unix_timestamp(col("date")))

# 定义窗口:按时间戳排序,范围为当前时间戳往前推2天(2*24*3600秒)到当前时间戳
window_spec = Window.orderBy(col("unix_date")) \
                   .rangeBetween(-2*24*3600, 0)

# 计算new_price并保留需要的列
result_df = df.withColumn("new_price", max(col("price")).over(window_spec)) \
              .select("date", "price", "new_price")

# 展示结果
result_df.show()

方法二:基于行的范围(仅适用于日期连续的场景)

如果确定数据源的日期是连续无缺失的,可以简化为按行范围计算:

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import col, max

spark = SparkSession.builder.appName("RowRangeMaxPrice").getOrCreate()

data = [
    ("1995-10-24", 10),
    ("1995-10-25", 20),
    ("1995-10-26", 30),
    ("1995-10-27", 11),
    ("1995-10-28", 14),
    ("1995-10-29", 12)
]
df = spark.createDataFrame(data, ["date", "price"])
df = df.withColumn("date", col("date").cast("date"))

# 定义窗口:按日期排序,范围为当前行往前2行到当前行
window_spec = Window.orderBy(col("date")) \
                   .rowsBetween(-2, 0)

result_df = df.withColumn("new_price", max(col("price")).over(window_spec))
result_df.show()

内容的提问来源于stack exchange,提问作者Binayak Rout

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 10:48:21