Оконные функции с rowsBetween в PySpark
Клауза ROWS BETWEEN: unboundedPreceding, unboundedFollowing, currentRow, относительные смещения. Накопительная сумма, скользящее среднее, исключение текущей строки.
Введение¶
Оконные функции в PySpark позволяют выполнять вычисления по строкам, связанным с текущей строкой, в пределах определённого окна. Клауза ROWS BETWEEN позволяет задать диапазон строк, которые оконная функция будет учитывать.
Синтаксис¶
from pyspark.sql.window import Window
# Определение оконной спецификации
window_spec = Window.partitionBy(<columns>).orderBy(<column>).rowsBetween(<start>, <end>)
# Применение оконной функции
df.withColumn("new_column", <window_function>("column").over(window_spec))
Ключевые концепции¶
partitionBy(<columns>): Разбивает данные на группы по указанным колонкамorderBy(<column>): Упорядочивает строки внутри каждой группыrowsBetween(<start>, <end>): Определяет диапазон строк для каждой строки
Параметры ROWS BETWEEN¶
Параметр start¶
Window.unboundedPreceding: Начало от первой строкиx PRECEDING(например,-2): Начало заxстрок до текущейWindow.currentRow: Начало с текущей строки
Параметр end¶
Window.unboundedFollowing: Конец на последней строкеx FOLLOWING(например,2): Конец черезxстрок после текущейWindow.currentRow: Конец на текущей строке
Примеры¶
1. Накопительная сумма¶
Вычисление нарастающего итога от первой строки до текущей:
from pyspark.sql.functions import sum
# Оконная спецификация для накопительной суммы
window_spec = Window.partitionBy("category").orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow)
# Применение накопительной суммы
df = df.withColumn("cumulative_sales", sum("sales").over(window_spec))
df.show()
Window.unboundedPreceding: Начало от первой строкиWindow.currentRow: Конец на текущей строке
2. Скользящее среднее (окно 3 дня)¶
Вычисление среднего за последние 3 дня, начиная за 2 строки до текущей:
from pyspark.sql.functions import avg
# Оконная спецификация для скользящего среднего по последним 3 строкам
window_spec = Window.partitionBy("category").orderBy("date").rowsBetween(-2, Window.currentRow)
# Применение скользящего среднего
df = df.withColumn("moving_avg_sales", avg("sales").over(window_spec))
df.show()
-2: Начало за 2 строки до текущейWindow.currentRow: Конец на текущей строке
3. Исключение текущей и будущих строк¶
Вычисление накопительной суммы без учёта текущей и будущих строк:
from pyspark.sql.functions import sum
# Оконная спецификация, исключающая текущую и будущие строки
window_spec = Window.partitionBy("category").orderBy("date").rowsBetween(Window.unboundedPreceding, -1)
# Применение накопительной суммы без текущей строки
df = df.withColumn("cumulative_sales_excluding_future", sum("sales").over(window_spec))
df.show()
Window.unboundedPreceding: Начало от первой строки-1: Конец на строке перед текущей
Полный пример¶
from pyspark.sql import SparkSession
from pyspark.sql.functions import sum, avg
from pyspark.sql.window import Window
# Тестовые данные
data = [
("2024-01-01", "A", 100),
("2024-01-02", "A", 150),
("2024-01-03", "A", 200),
("2024-01-01", "B", 50),
("2024-01-02", "B", 75),
("2024-01-03", "B", 100)
]
# Создание DataFrame
spark = SparkSession.builder.appName("WindowFunctions").getOrCreate()
df = spark.createDataFrame(data, ["date", "category", "sales"])
# 1. Накопительная сумма
window_spec = Window.partitionBy("category").orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow)
df_cumulative_sum = df.withColumn("cumulative_sales", sum("sales").over(window_spec))
# 2. Скользящее среднее (окно 3 дня)
window_spec_avg = Window.partitionBy("category").orderBy("date").rowsBetween(-2, Window.currentRow)
df_moving_avg = df_cumulative_sum.withColumn("moving_avg_sales", avg("sales").over(window_spec_avg))
# 3. Исключение текущей и будущих строк
window_spec_excluding_future = Window.partitionBy("category").orderBy("date").rowsBetween(Window.unboundedPreceding, -1)
df_final = df_moving_avg.withColumn("cumulative_sales_excluding_current_n_future", sum("sales").over(window_spec_excluding_future))
# Показать результат
df_final.show()
Ожидаемый вывод:
+----------+--------+-----+-----------------+-----------------+-------------------------------------------+
| date|category|sales| cumulative_sales| moving_avg_sales|cumulative_sales_excluding_current_n_future|
+----------+--------+-----+-----------------+-----------------+-------------------------------------------+
|2024-01-01| A| 100| 100| 100.0| null|
|2024-01-02| A| 150| 250| 125.0| 100|
|2024-01-03| A| 200| 450| 150.0| 250|
|2024-01-01| B| 50| 50| 50.0| null|
|2024-01-02| B| 75| 125| 62.5| 50|
|2024-01-03| B| 100| 225| 75.0| 125|
+----------+--------+-----+-----------------+-----------------+-------------------------------------------+
Итог¶
Клауза ROWS BETWEEN позволяет задать конкретные диапазоны строк вокруг текущей строки для аналитических задач: накопительных сумм, скользящих средних и других вычислений. Управляя параметрами start и end, можно гибко и точно настраивать сложные вычисления над данными.