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

PySpark:按5分钟窗口将多列合并为JSON结构列

解决方案

核心思路

通过宽表转长表统一处理所有Alarm列,结合窗口函数计算告警持续时长,最后聚合生成目标JSON结构,天然支持最多10个Alarm列的扩展需求。


具体实现步骤

1. 计算单条记录的有效持续时长

按system分组、time升序排序,用lead函数获取下一条记录的时间,计算当前记录到下一条的时间差(秒),最后一条记录的持续时长设为0(可根据实际需求调整)。

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 定义排序窗口:按系统分组,时间升序
time_sort_window = Window.partitionBy("system").orderBy("time")

# 计算每条记录的后续时间及持续秒数
df_with_duration = df.withColumn(
    "next_time", F.lead("time").over(time_sort_window)
).withColumn(
    "duration_sec",
    F.coalesce(F.unix_timestamp("next_time") - F.unix_timestamp("time"), F.lit(0))
)

2. 将多列Alarm转换为长表结构

用stack函数把Alarm1到Alarm10批量转成alarm_value列,同时过滤空告警值,并去重同一时间戳下的重复告警值(避免重复计算)。

# 构造stack表达式,支持10个Alarm列
num_alarm_cols = 10
stack_expr = f"stack({num_alarm_cols}, " + ", ".join([f"'Alarm{i}', Alarm{i}" for i in range(1, num_alarm_cols+1)]) + ")"

# 宽表转长表并去重
df_long = df_with_duration.select(
    "system", "time", "duration_sec",
    F.explode(F.stack(stack_expr)).alias("alarm_col", "alarm_value")
).filter(
    F.col("alarm_value").isNotNull()
).dropDuplicates(["system", "time", "alarm_value"])

3. 5分钟窗口聚合告警时长

按system和5分钟对齐的窗口分组,计算每个告警值在窗口内的总持续时长。

# 定义5分钟滚动窗口:按系统分组,窗口起始时间对齐到5分钟边界
window_agg_spec = Window.partitionBy(
    "system",
    F.date_trunc("5 minutes", "time").alias("window_start")
)

# 聚合每个窗口内的告警总时长
df_window_sum = df_long.withColumn(
    "total_duration",
    F.sum("duration_sec").over(window_agg_spec.partitionBy("system", "window_start", "alarm_value"))
).select(
    "system",
    F.col("window_start").alias("start_time"),
    F.date_add(F.col("window_start"), 5*60).alias("end_time"),
    "alarm_value",
    F.first("total_duration").over(window_agg_spec.partitionBy("system", "window_start", "alarm_value")).alias("total_duration")
).dropDuplicates(["system", "start_time", "alarm_value"])

4. 生成JSON格式的Alarms列

用map_from_arrays将告警值和时长映射为键值对,再转成JSON字符串。

# 分组生成最终结果
final_df = df_window_sum.groupBy(
    "system", "start_time", "end_time"
).agg(
    F.to_json(
        F.map_from_arrays(
            F.collect_list("alarm_value"),
            F.collect_list("total_duration")
        )
    ).alias("Alarms")
)

关键说明

  • 扩展性:只需调整num_alarm_cols变量为10,即可直接支持Alarm1到Alarm10,无需修改核心逻辑。
  • 去重逻辑:通过dropDuplicates(["system", "time", "alarm_value"])确保同一时间戳下的重复告警值不会被重复计算,符合需求要求。
  • 窗口对齐:用date_trunc实现严格的5分钟滚动窗口,若需要滑动窗口,可改用window函数定义滑动规则。

内容的提问来源于stack exchange,提问作者user23139627

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 06:54:53