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
相关产品推荐
相关产品推荐

