基于年龄与故障状态使用PySpark生成子样本的需求
样本采集需求
核心规则
- 采集时长范围:3天
- 排除项:健康状态(failure_status为No)且使用时长不足3天的序列号,无需纳入
- 纳入项:
- 故障状态(failure_status为Yes)的序列号:需保留故障发生前≤3天的所有数据
- 健康状态且使用时长≥3天的序列号:需保留最近3天的所有数据
示例说明
- 序列号C于1月3日故障,纳入其1月1日、2日的数据
- 序列号D于1月4日故障,纳入其1月1日、2日、3日的数据
- 序列号A、B为健康状态,纳入其1月3日、4日、5日的最近3天数据
- 序列号E、F为健康且使用时长不足3天,无需纳入
初始加载代码
url="https://gist.githubusercontent.com/JishanAhmed2019/6625009b71ade22493c256e77e1fdaf3/raw/8b51625b76a06f7d5c76b81a116ded8f9f790820/FailureSample.csv" from pyspark import SparkFiles spark.sparkContext.addFile(url) df=spark.read.csv(SparkFiles.get("FailureSample.csv"), header=True,sep='\t')
数据格式说明
当前数据格式

预期样本格式

实现代码
通过Spark SQL的窗口函数和日期函数完成规则筛选:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 转换date列为日期类型 df = df.withColumn("date", F.to_date("date", "yyyy-MM-dd")) # 按序列号分组的窗口定义 serial_window = Window.partitionBy("serial_number") # 计算每个序列号的关键指标:最近日期、故障日期、总使用天数 df_enhanced = df.withColumn("latest_date", F.max("date").over(serial_window)) \ .withColumn("fail_date", F.when(F.col("failure_status") == "Yes", F.col("date")).over(serial_window)) \ .withColumn("fail_date", F.max("fail_date").over(serial_window)) \ .withColumn("total_usage_days", F.datediff(F.col("latest_date"), F.min("date").over(serial_window)) + 1) # 应用筛选规则 final_df = df_enhanced.filter( # 故障样本:日期在故障日期前3天内(不包含故障当天) (F.col("failure_status") == "Yes") & (F.datediff(F.col("fail_date"), F.col("date")) <= 3) & (F.col("date") < F.col("fail_date")) | # 健康样本:使用时长≥3天,且日期属于最近3天 (F.col("failure_status") == "No") & (F.col("total_usage_days") >= 3) & (F.datediff(F.col("latest_date"), F.col("date")) <= 2) ).drop("latest_date", "fail_date", "total_usage_days") # 查看结果 final_df.orderBy("serial_number", "date").show()
内容的提问来源于stack exchange,提问作者ForestGump
相关产品推荐
相关产品推荐

