在PySpark DataFrame中新增列:计算指定日期区间价格最大值
PySpark实现日期区间内价格最大值计算
现有DataFrame
| date | price |
|---|---|
| 1995-10-24 | 10 |
| 1995-10-25 | 20 |
| 1995-10-26 | 30 |
| 1995-10-27 | 11 |
| 1995-10-28 | 14 |
| 1995-10-29 | 12 |
需求说明
新增一列new_price,取值为当前日期与该日期前2天区间内的price最大值,即new_price = max(price) between (date, date-2 days)。预期结果如下:
预期结果
| date | price | new_price |
|---|---|---|
| 1995-10-24 | 10 | 10 |
| 1995-10-25 | 20 | 20 |
| 1995-10-26 | 30 | 30 |
| 1995-10-27 | 11 | 30 |
| 1995-10-28 | 14 | 30 |
| 1995-10-29 | 12 | 14 |
代码实现
方法一:基于时间戳的区间范围(推荐,适配日期不连续场景)
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
相关产品推荐
相关产品推荐

