这是一种通过确定每个值的行号来完成此操作的方法。
使用pyspark.sql.functions.monotonically_increasing_id() 创建一个唯一的、递增的id 列。
创建一个在id 列上执行orderBy() 的pyspark.sql.Window()。
在窗口上方使用pyspark.sql.functions.row_number() 来获取每个值的行号。
将 row_number - 1(因为它从 1 开始)除以组数,然后发言以获得组数。
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|
#+-----+----+