【问题标题】:Apply a function every 60 rows in a pyspark dataframe在 pyspark 数据框中每 60 行应用一个函数
【发布时间】: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


【解决方案1】:

我认为groupBy 足以获得所需的结果。

df.show()
+---+------+------+
| HR|maxABP|Second|
+---+------+------+
|110| 128.0|    10|
|110| 127.0|    20|
|111| 127.0|    30|
|111| 127.0|    40|
|111| 126.0|    50|
|111| 127.0|    60|
|109| 126.0|    70|
|111| 126.0|    80|
+---+------+------+

df.withColumn('Minute', f.expr('cast(Second / 60 as int)')) \
  .groupBy('Minute').agg( \
    f.round(f.min('HR'), 2).alias('Min_HR'), \
    f.round(f.max('HR'), 2).alias('Max_HR'), \
    f.round(f.avg('HR'), 2).alias('Avg_HR'), \
    f.max('maxABP').alias('maxABP')) \
  .withColumn('Alarm', f.expr('if(maxABP < 85, 1, 0)')) \
  .show()

+------+------+------+------+------+-----+
|Minute|Min_HR|Max_HR|Avg_HR|maxABP|Alarm|
+------+------+------+------+------+-----+
|     1|   109|   111|110.33| 127.0|    0|
|     0|   110|   111| 110.6| 128.0|    0|
+------+------+------+------+------+-----+

【讨论】:

  • 是的,这就是我最终解决它的方式,感谢您的回复。我通过使用获得了分钟列: withColumn('Minute', F.ceil(df_filtered['Second'] / 60))
【解决方案2】:

示例 DF:

df = spark.createDataFrame(
  [
     (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, 1001)
    ,(114, 126.0, 1003),(115, 83.0, 1064),(116, 127.0, 1066)
  ], ['HR', 'maxABP', 'Second']
)

+---+------+------+
| 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|  1001|
|114| 126.0|  1003|
|115|  83.0|  1064|
|116| 127.0|  1066|

然后使用窗口函数:

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

w1 = (Window.partitionBy(F.col('Minute')))

df\
  .withColumn('Minute', F.round(F.col('Second')/60,0)+1)\
  .withColumn('Min_HR', F.min('HR').over(w1))\
  .withColumn('Max_HR', F.max('HR').over(w1))\
  .withColumn('Avg_HR', F.round(F.avg('HR').over(w1),0))\
  .withColumn('Min_ABP', F.round(F.min('maxABP').over(w1),0))\
  .select('Min_HR','Max_HR','Min_ABP','Avg_HR','Minute')\
  .dropDuplicates()\
  .withColumn('Alarm', F.when(F.col('Min_ABP')<85, 1).otherwise(F.lit('0')))\
  .select('Min_HR','Max_HR','Avg_HR','Alarm','Minute')\
  .orderBy('Minute')\
  .show()

+------+------+------+-----+------+
|Min_HR|Max_HR|Avg_HR|Alarm|Minute|
+------+------+------+-----+------+
|   109|   111| 110.0|    0|   1.0|
|   111|   114| 113.0|    0|  18.0|
|   115|   116| 116.0|    1|  19.0|

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-10-29
    • 2011-04-08
    • 2011-06-15
    • 1970-01-01
    • 2011-11-02
    • 1970-01-01
    • 2017-06-15
    • 1970-01-01
    相关资源
    最近更新 更多