PySpark使用窗口函数计算指定日期范围内故障点倒计数RUL
需求说明
现有设备运行数据表结构如下:
| equipment | run | runend | failure | removal_date |
|---|---|---|---|---|
| A | 1 | 1/1/2021 | 0 | 4/1/2021 |
| A | 2 | 2/1/2021 | 0 | 4/1/2021 |
| A | 3 | 3/1/2021 | 0 | 4/1/2021 |
| A | 4 | 4/1/2021 | 1 | 4/1/2021 |
| A | 5 | 4/1/2021 | 0 | 20/1/2021 |
| A | 6 | 10/1/2021 | 0 | 20/1/2021 |
需要新增RUL列,作为到故障终止点的倒计数,终止点为同一设备同一次移除周期内,和removal_date最接近的runend值,最终输出效果如下:
| equipment | run | runend | failure | removal_date | RUL |
|---|---|---|---|---|---|
| A | 1 | 1/1/2021 | 0 | 4/1/2021 | 3 |
| A | 2 | 2/1/2021 | 0 | 4/1/2021 | 2 |
| A | 3 | 3/1/2021 | 0 | 4/1/2021 | 1 |
| A | 4 | 4/1/2021 | 1 | 4/1/2021 | 0 |
| A | 5 | 4/1/2021 | 0 | 20/1/2021 | 16 |
| A | 6 | 10/1/2021 | 0 | 20/1/2021 | 10 |
现有问题
当前编写的窗口函数分区逻辑错误,按equipment和run两个字段分区,每个run为独立分组,无法实现跨run的周期内倒计数需求:
w = Window.partitionBy("equipment", "run").orderBy(asc("runend")) df = df.withColumn("rank", rank().over(w)) # 查看DataFrame结构 df.where(col("equipment") == "A").groupby("equipment", "run", "rank", "failure", "runend", "removal_date").count().orderBy("equipment", "runend").show()
运行后得到的rank是全局连续排序,不是按移除周期的倒计数:
| equipment | run | runend | failure | removal_date | rank |
|---|---|---|---|---|---|
| A | 1 | 1/1/2021 | 0 | 4/1/2021 | 1 |
| A | 2 | 2/1/2021 | 0 | 4/1/2021 | 2 |
| A | 3 | 3/1/2021 | 0 | 4/1/2021 | 3 |
| A | 4 | 4/1/2021 | 1 | 4/1/2021 | 4 |
| A | 5 | 4/1/2021 | 0 | 20/1/2021 | 5 |
| A | 6 | 10/1/2021 | 0 | 20/1/2021 | 6 |
解决方案
核心逻辑调整
- 首先将日期字段从字符串转为日期类型,避免字符串计算偏差
- 窗口分区改为按
equipment和removal_date分组,对应同一个设备的同一次移除周期 - 取每个分组内最大的
runend作为故障终止点,计算当前行runend到终止点的天数差即为RUL
完整代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 转换日期格式,注意输入日期为日/月/年格式 df = df.withColumn("runend", F.to_date(F.col("runend"), "d/M/yyyy")) \ .withColumn("removal_date", F.to_date(F.col("removal_date"), "d/M/yyyy")) # 2. 定义同设备同移除周期的窗口 w = Window.partitionBy("equipment", "removal_date") # 3. 计算每个周期的终止runend,再算日期差得到RUL df = df.withColumn("cycle_end_runend", F.max("runend").over(w)) \ .withColumn("RUL", F.datediff(F.col("cycle_end_runend"), F.col("runend"))) \ .drop("cycle_end_runend") # 测试查看结果 df.orderBy("equipment", "runend").show()
输出验证
运行上述代码后得到的结果和预期完全一致。
内容的提问来源于stack exchange,提问作者Aesir
相关产品推荐
相关产品推荐

