PySpark DataFrame每60行分组聚合计算心率指标与maxABP告警
PySpark 按每60行/60秒分组聚合实现方案
核心思路:首先给每行分配所属的分钟分组标签,再基于分组做聚合计算即可,不需要手动实现原生MapReduce逻辑,PySpark内置算子已经封装了分布式计算能力。
场景1:Second字段为连续秒数,按实际秒数归为1分钟
直接通过Second字段计算分组即可:
from pyspark.sql import functions as F # 1. 计算每行所属的分钟分组:1-60秒为第1分钟,61-120为第2分钟,以此类推 df_with_minute = df.withColumn("Minute", ((F.col("Second") - 1) // 60).cast("int") + 1) # 2. 分组聚合计算所需指标 result_df = df_with_minute.groupBy("Minute").agg( F.min("HR").alias("Min_HR"), F.max("HR").alias("Max_HR"), F.avg("HR").alias("Avg_HR"), # 只要组内有任意一个maxABP<85,Alarm就为1,否则为0 F.max(F.when(F.col("maxABP") < 85, 1).otherwise(0)).alias("Alarm") ).orderBy("Minute") # 查看结果 result_df.show()
场景2:严格按每60行分组(不管Second是否连续)
如果需要严格每60行分一组,和Second值无关,就先给每行生成全局连续行号再计算分组:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 按Second排序生成全局行号(如果有其他排序逻辑可以调整orderBy的字段) window_spec = Window.orderBy("Second") df_with_row = df.withColumn("row_num", F.row_number().over(window_spec)) # 计算所属分钟分组 df_with_minute = df_with_row.withColumn("Minute", ((F.col("row_num") - 1) // 60).cast("int") + 1) # 后续聚合逻辑和场景1完全一致 result_df = df_with_minute.groupBy("Minute").agg( F.min("HR").alias("Min_HR"), F.max("HR").alias("Max_HR"), F.avg("HR").alias("Avg_HR"), F.max(F.when(F.col("maxABP") < 85, 1).otherwise(0)).alias("Alarm") ).orderBy("Minute") result_df.show()
原有代码错误说明
- PySpark DataFrame是分布式存储的数据集,不支持Python本地列表的
df[x:y]切片操作,该写法只能作用于本地Python列表。 - 你写的
map/reduce是Python本地的函数,不能直接作用于分布式DataFrame,需要使用PySpark提供的分布式算子。 - 手动实现原生MapReduce逻辑完全可以满足需求,但PySpark内置的分组聚合算子已经做了大量优化,性能远高于自行实现的MapReduce逻辑,不需要重复造轮子。
内容的提问来源于stack exchange,提问作者Bijlenman
相关产品推荐
相关产品推荐

