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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 04:24:05