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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 19:05:22