【问题标题】:Average over 2000 values with PySpark Dataframe使用 PySpark Dataframe 平均超过 2000 个值
【发布时间】:2018-08-23 17:43:48
【问题描述】:

我有一个大约十亿行的 PySpark 数据框。我想对每 2000 个值进行平均,例如索引为 0-1999 的行的平均值,索引为 2000-3999 的行的平均值,等等。我该怎么做呢?或者,我也可以为每 2000 个平均 10 个值,例如索引为 0-9 的行的平均值,索引为 2000-2009 的行的平均值,等等。这样做的目的是对数据进行下采样。我目前没有索引行,所以如果我需要这个,我该怎么做?

【问题讨论】:

  • 你必须对数据进行分区,将2000行划分为一个分区,然后使用mapPartition api获取平均值

标签: python apache-spark pyspark bigdata spark-dataframe


【解决方案1】:

您可以使用monotonically_increasing_id() 生成行ID,将其划分并使用上限函数在您想要的任何间隔内生成ID。然后使用窗口函数对该 id 进行分区并生成平均值。例如,假设您的数据框是 data 并且您希望对列 value 进行平均,则应该可以使用以下内容。

import org.apache.spark.sql.expressions.Window
val partitionWindow = Window.partitionBy($"rowId")
data.withColumn("rowId", floor(monotonically_increasing_id()/2000.0)).withColumn("avg", avg(data("value")) over(partitionWindow)).show()

希望对您有所帮助。

【讨论】:

  • 我认为这不会起作用,因为monotonically_increasing_id() 保证会增加但不一定连续。因此,第 1 行的 id 值可能为 0,第 2 行的 id 值可能为 1000000。
  • @pault 呵呵,我在我正在使用的几个数据帧上进行了尝试,它似乎生成了连续的数字。但是经过一番搜索,该功能似乎确实存在很多问题。我认为,来自 spark-sql 的 row_number() 应该可以解决问题。
  • 你可能需要两者的组合,但我想它会很慢。
【解决方案2】:

这是一种通过确定每个值的行号来完成此操作的方法。

  1. 使用pyspark.sql.functions.monotonically_increasing_id() 创建一个唯一的、递增的id 列。

  2. 创建一个在id 列上执行orderBy()pyspark.sql.Window()

  3. 在窗口上方使用pyspark.sql.functions.row_number() 来获取每个值的行号。

  4. 将 row_number - 1(因为它从 1 开始)除以组数,然后发言以获得组数。

  5. groupBy() 组数并计算平均值。


这是一个例子:

创建示例数据

对于这个例子,我将创建一个包含 5 个连续值的数据框,从 10 到 40(包括在内)的每个 10 的倍数开始。此示例中的组大小为 5 - 我们需要 5 个连续值的平均值。

data = map(
    lambda y: (y, ),
    reduce(
        list.__add__,
        [range(x, x+5) for x in range(10, 50, 10)]
    )
)
df = sqlCtx.createDataFrame(data, ["col1"])
df.show()
#+----+
#|col1|
#+----+
#|  10|
#|  11|
#|  12|
#|  13|
#|  14|
#|  20|
#|  21|
#|  22|
#|  23|
#|  24|
#|  30|
#|  31|
#|  32|
#|  33|
#|  34|
#|  40|
#|  41|
#|  42|
#|  43|
#|  44|
#+----+

添加 ID 列

我展示这个步骤是为了证明monotonically_increasing_id() 不能保证是连续的。

import pyspark.sql.functions as f
df = df.withColumn('id', f.monotonically_increasing_id())
df.show()
#+----+----------+
#|col1|        id|
#+----+----------+
#|  10|         0|
#|  11|         1|
#|  12|         2|
#|  13|         3|
#|  14|         4|
#|  20|         5|
#|  21|         6|
#|  22|         7|
#|  23|         8|
#|  24|         9|
#|  30|8589934592|
#|  31|8589934593|
#|  32|8589934594|
#|  33|8589934595|
#|  34|8589934596|
#|  40|8589934597|
#|  41|8589934598|
#|  42|8589934599|
#|  43|8589934600|
#|  44|8589934601|
#+----+----------+

计算组数

from pyspark.sql import Window
group_size = 5
w = Window.orderBy('id')
df = df.withColumn('group', f.floor((f.row_number().over(w) - 1) / group_size))\
    .select('col1', 'group')
df.show()
#+----+-----+
#|col1|group|
#+----+-----+
#|  10|    0|
#|  11|    0|
#|  12|    0|
#|  13|    0|
#|  14|    0|
#|  20|    1|
#|  21|    1|
#|  22|    1|
#|  23|    1|
#|  24|    1|
#|  30|    2|
#|  31|    2|
#|  32|    2|
#|  33|    2|
#|  34|    2|
#|  40|    3|
#|  41|    3|
#|  42|    3|
#|  43|    3|
#|  44|    3|
#+----+-----+

获取每组的平均值

df.groupBy('group').agg(f.avg('col1').alias('avg')).show()
#+-----+----+
#|group| avg|
#+-----+----+
#|    0|12.0|
#|    1|22.0|
#|    2|32.0|
#|    3|42.0|
#+-----+----+

【讨论】:

    猜你喜欢
    • 2017-11-07
    • 1970-01-01
    • 2017-07-05
    • 2019-03-07
    • 1970-01-01
    • 2014-09-23
    • 1970-01-01
    • 2021-05-05
    • 2018-07-11
    相关资源
    最近更新 更多