Оконные функции с rowsBetween в PySpark

Клауза ROWS BETWEEN: unboundedPreceding, unboundedFollowing, currentRow, относительные смещения. Накопительная сумма, скользящее среднее, исключение текущей строки.

core

Введение

Оконные функции в 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, можно гибко и точно настраивать сложные вычисления над данными.