使用PySpark计算系统状态降级的持续时间
PySpark 实现系统降级时间段统计方案
核心思路
通过窗口函数和状态分组,将连续的非green状态与后续第一个green状态关联,形成完整的降级周期,再统计每个周期的关键信息。
步骤实现
1. 初始化数据与环境
先将示例数据加载为PySpark DataFrame,并确保timestamp列转换为时间类型:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("DegradationAnalysis").getOrCreate() # 示例数据 data = [ ("red", "2023-01-02T01:05:32.113Z"), ("yellow", "2023-01-02T01:15:47.329Z"), ("red", "2023-01-02T01:25:11.257Z"), ("green", "2023-01-02T01:33:12.187Z"), ("red", "2023-01-05T15:10:12.854Z"), ("green", "2023-01-05T15:26:24.131Z") ] # 创建DataFrame并转换timestamp类型 df = spark.createDataFrame(data, ["status", "timestamp"]) df = df.withColumn("timestamp", F.to_timestamp("timestamp"))
2. 标记降级分组
按时间排序,用窗口函数标记每个状态是否属于新的降级周期,再通过累积求和生成唯一的降级组ID:
# 定义窗口:按时间升序排序 window = Window.orderBy("timestamp") # 标记新降级组的起始点 df = df.withColumn( "is_new_degradation", F.when( (F.col("status") != "green") & (F.lag("status").over(window) == "green"), 1 ).when( (F.col("status") != "green") & (F.lag("status").over(window).isNull()), 1 ).otherwise(0) ) # 生成降级组ID df = df.withColumn( "degradation_group", F.sum("is_new_degradation").over(window.rangeBetween(Window.unboundedPreceding, 0)) )
3. 提取降级周期的起止时间与时长
对每个降级组,筛选非green状态的最早时间作为开始,对应green状态的时间作为结束,计算持续时长:
# 分组聚合获取每个降级组的开始时间 start_df = df.filter(F.col("status") != "green") \ .groupBy("degradation_group") \ .agg(F.min("timestamp").alias("start")) # 分组聚合获取每个降级组的结束时间(对应green状态的时间) end_df = df.filter(F.col("status") == "green") \ .withColumn( "degradation_group", F.lag("degradation_group").over(window) ) \ .groupBy("degradation_group") \ .agg(F.min("timestamp").alias("end")) # 关联起止时间,计算持续时长(分钟,保留两位小数) result_df = start_df.join(end_df, on="degradation_group", how="inner") \ .withColumn( "duration", F.round(F.minute(F.col("end") - F.col("start")) + F.second(F.col("end") - F.col("start"))/60, 2) ) \ .withColumn( "degradation", F.concat(F.lit("degradation "), F.col("degradation_group").cast("string")) ) \ .select("degradation", "duration", "start", "end")
4. 查看结果
result_df.show(truncate=False)
输出与预期一致:
+-------------+--------+------------------------+------------------------+ |degradation |duration|start |end | +-------------+--------+------------------------+------------------------+ |degradation 1|27.65 |2023-01-02 01:05:32.113 |2023-01-02 01:33:12.187 | |degradation 2|16.2 |2023-01-05 15:10:12.854 |2023-01-05 15:26:24.131 | +-------------+--------+------------------------+------------------------+
关键说明
- 窗口函数按时间顺序追踪状态变化,确保降级组划分准确
- 分组聚合精准提取每个周期的起止点,排除中间状态干扰
- 时长计算通过
minute和second函数转换为分钟数,匹配预期输出格式
内容的提问来源于stack exchange,提问作者Connor
相关产品推荐
相关产品推荐

