使用PySpark重标记故障前2天的健康样本并修正Serial C数据问题
PySpark 故障数据修正方案
需求说明
- 针对所有序列号,将实际故障发生前2天的健康样本(标记为0)重标记为故障(标记为1)
- 修正Serial C在实际故障后仍被标记为健康的错误
完整实现代码
import findspark findspark.init() import pyspark from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化Spark会话 spark = SparkSession.builder.appName('failure_data_correction').getOrCreate() # 读取数据 url="https://gist.githubusercontent.com/JishanAhmed2019/e464ca4da5c871428ca9ed9264467aa0/raw/da3921c1953fefbc66dddc3ce238dac53142dba8/failure.csv" from pyspark import SparkFiles spark.sparkContext.addFile(url) df = spark.read.csv(SparkFiles.get("failure.csv"), header=True, sep='\t') # 转换日期格式并按序列号+日期排序 df = df.withColumn("date", F.to_date(F.col("date"), "yyyy-MM-dd")) \ .orderBy("serial_number", "date") # 提取每个序列号的故障日期 failure_dates = df.filter(F.col("failure") == 1) \ .select("serial_number", F.col("date").alias("failure_date")) # 关联原数据,计算每条记录距离故障日期的天数 df_with_failure_info = df.join(failure_dates, on="serial_number", how="left") \ .withColumn("days_before_failure", F.datediff(F.col("failure_date"), F.col("date"))) # 修正failure列标记 corrected_df = df_with_failure_info.withColumn( "failure", F.when( (F.col("days_before_failure").between(0, 2)) & (F.col("failure") == 0), 1 ).when( (F.col("date") > F.col("failure_date")) & (F.col("failure") == 0), 1 ).otherwise(F.col("failure")) ).drop("failure_date", "days_before_failure") # 查看修正结果 corrected_df.show()
关键逻辑说明
- 日期标准化:将字符串格式的
date转为日期类型,保证时间计算的准确性 - 故障日期定位:筛选出所有故障标记为1的记录,得到每个设备的故障时间点
- 时间差计算:通过关联操作,计算每条样本距离对应设备故障日期的天数
- 标记修正规则:
- 故障前0-2天的健康样本(标记0)改为1
- 故障日期之后的健康样本(标记0)改为1
- 其余情况保留原标记
内容的提问来源于stack exchange,提问作者ForestGump
相关产品推荐
相关产品推荐

