PySpark:如何用窗口函数计算指定日期范围的产品最低价格
解决PySpark按日期范围计算最低价格的问题
问题根源
你之前的窗口函数用了rowsBetween(30, Window.currentRow),这是按行的数量划定窗口范围,而非日期范围。你的数据集只有4行,当前行之前根本没有30行,所以计算结果返回null。而且这种方式无法保证窗口内的日期都在指定的N天范围内,完全不符合需求。
正确实现步骤
要实现按日期范围(30/60/90天)计算最低价格,需要基于日期的时间戳数值定义窗口范围,同时避免过滤和Join操作,具体步骤如下:
- 转换日期字段为Spark日期类型:原始
Date是字符串,无法直接进行日期计算,先转换成标准日期类型。 - 基于时间戳定义窗口范围:将日期转换为Unix时间戳(秒数),利用
rangeBetween划定当前日期往前N天的时间范围,确保窗口内仅包含符合日期要求的数据。
完整代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 转换Date列为Spark日期类型(匹配原始数据的MM/dd/yyyy格式) df = df.withColumn("Date", F.to_date("Date", "MM/dd/yyyy")) # 2. 定义基础窗口:按产品类型分区,按日期的Unix时间戳排序 base_window = Window.partitionBy("Product_Type").orderBy(F.unix_timestamp("Date")) # 3. 分别计算最近30、60、90天的最低价格 # 30天对应秒数:30*24*3600 = 2592000 df = df.withColumn("minPrice30", F.min("Price").over(base_window.rangeBetween(-2592000, 0))) # 60天对应秒数:60*24*3600 = 5184000 df = df.withColumn("minPrice60", F.min("Price").over(base_window.rangeBetween(-5184000, 0))) # 90天对应秒数:90*24*3600 = 7776000 df = df.withColumn("minPrice90", F.min("Price").over(base_window.rangeBetween(-7776000, 0))) # 调整列顺序匹配预期输出 df = df.select("minPrice60", "minPrice30", "Price", "Date", "InvoiceNumber", "Product_Type") df.show()
结果验证
运行上述代码后,会得到与你预期一致的输出:
| minPrice60 | minPrice30 | Price | Date | InvoiceNumber | Product_Type |
|---|---|---|---|---|---|
| 12.65 | 12.65 | 12.65 | 2021-12-30 | INV_19984 | AXN UN1234 |
| 12.65 | 12.65 | 18.78 | 2022-01-23 | INV_200174 | AXN UN1234 |
| 11.78 | 11.78 | 11.78 | 2022-01-25 | INV_200173 | AXN UN1234 |
| 11.1 | 11.1 | 11.1 | 2022-03-02 | INV_9912 | AXN UN1234 |
关键说明
- rangeBetween vs rowsBetween:
rangeBetween基于排序字段的数值范围划定窗口,这里用Unix时间戳差值精准控制日期范围;rowsBetween基于行的数量,完全不适合日期范围需求。 - 性能优化:这种方式仅通过窗口函数一次遍历数据完成计算,不需要过滤或Join操作,能高效处理大规模数据集,符合你的性能要求。
内容的提问来源于stack exchange,提问作者Ramaraju.d
相关产品推荐
相关产品推荐

