【发布时间】:2022-01-07 22:04:06
【问题描述】:
我的数据框名为 df,有 123729 行,如下所示:
+---+------+------+
| HR|maxABP|Second|
+---+------+------+
|110| 128.0| 1|
|110| 127.0| 2|
|111| 127.0| 3|
|111| 127.0| 4|
|111| 126.0| 5|
|111| 127.0| 6|
|109| 126.0| 7|
|111| 126.0| 8|
我需要每 60 行(或秒)聚合多个值。对于每一分钟,我都想知道最小心率、平均心率、最大心率,以及在这些秒内的任何一秒内 maxABP 是否低于 85。所需的输出如下表所示,如果 maxABP
| Min_HR | Max_HR | Avg_HR | Alarm | Minute |
|---|---|---|---|---|
| 70 | 100 | 80 | 1 | 1 |
| 60 | 90 | 75 | 0 | 2 |
我想知道是否可以使用 mapreduce 将每 60 行聚合为这些单个值。我知道有很多错误,但可能是这样的:
def max_HR(df, i):
x = i
y = i+60
return reduce(lambda x, y: max(df[x:y]))
df_maxHR = map(lambda i: max_HR(i))
i 应该是 df 的一部分。
【问题讨论】:
-
添加更多示例数据。是否有分钟列,或者第二列超过 60,例如 61、62 等等?
-
所有样本数据都是这样的。没有更多的列,秒数继续达到 123729
-
另外,这个 DF 是如何填充的?通过流媒体?总是固定数量的 123729 行?
-
在并行化数据时是否可以在 60 行分区中这样做而不是使用默认值?
sc.parallelize(<my_data>, <num_rows>%60)
标签: python dataframe pyspark mapreduce