You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 20:10:32