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

Window specification: partitionBy и orderBy. Агрегатные функции в оконном контексте (sum, avg, min, max, count). Ранжирующие функции: rank, dense_rank, row_number.

core

Window Specification

Для определения окна в PySpark используется спецификация Window, которая включает:

  1. partitionBy: Разбивает данные на группы по одной или нескольким колонкам.
  2. orderBy: Упорядочивает строки внутри каждой партиции.
from pyspark.sql import Window

window_spec = Window.partitionBy("category").orderBy("sales")

Агрегатные функции в оконном контексте

Агрегатные функции sum, avg, min, max и другие применяются к определённому окну и позволяют вычислять нарастающие итоги, средние значения и другие показатели.

Функция Описание
sum Нарастающий итог по окну
avg Скользящее среднее
min Минимальное значение
max Максимальное значение
count Количество строк в окне

Примеры использования

Нарастающий итог:

from pyspark.sql.functions import col, sum

df = df.withColumn("cumulative_sales", sum("sales").over(window_spec))

Максимальные продажи в партиции:

from pyspark.sql.functions import max

df = df.withColumn("max_sales", max("sales").over(window_spec))

Среднее значение продаж:

from pyspark.sql.functions import avg

df = df.withColumn("average_sales", avg("sales").over(window_spec))

Ранжирующие функции

Ранжирующие функции присваивают ранг или порядковый номер строкам внутри оконной партиции.

Функция Описание Пример вывода
rank Ранг с пропуском при совпадении 1, 2, 2, 4
dense_rank Ранг без пропусков 1, 2, 2, 3
row_number Уникальный порядковый номер 1, 2, 3, 4

Ранжирование товаров по продажам:

from pyspark.sql.functions import rank

df = df.withColumn("rank", rank().over(window_spec))

Dense rank для лидеров:

from pyspark.sql.functions import dense_rank

df = df.withColumn("dense_rank", dense_rank().over(window_spec))

Уникальные порядковые номера:

from pyspark.sql.functions import row_number

df = df.withColumn("row_number", row_number().over(window_spec))

Комбинирование нескольких оконных функций

from pyspark.sql.functions import sum, rank, avg

df = df \
    .withColumn("rank", rank().over(window_spec)) \
    .withColumn("cumulative_sales", sum("sales").over(window_spec)) \
    .withColumn("average_sales", avg("sales").over(window_spec))

Полный пример

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum, avg, rank, dense_rank, row_number
from pyspark.sql.window import Window

# Тестовые данные
data = [
    ("Electronics", "Phone", 1000),
    ("Electronics", "Laptop", 1500),
    ("Electronics", "Tablet", 800),
    ("Furniture", "Chair", 300),
    ("Furniture", "Table", 300),
    ("Furniture", "Desk", 600),
]

# Создание DataFrame
spark = SparkSession.builder.appName("WindowFunctions").getOrCreate()
df = spark.createDataFrame(data, ["category", "product", "sales"])

# Определение оконной спецификации
window_spec = Window.partitionBy("category").orderBy("sales")

# Применение оконных функций
df_transformed = df \
    .withColumn("rank", rank().over(window_spec)) \
    .withColumn("dense_rank", dense_rank().over(window_spec)) \
    .withColumn("row_number", row_number().over(window_spec)) \
    .withColumn("cumulative_sales", sum("sales").over(window_spec)) \
    .withColumn("average_sales", avg("sales").over(window_spec))

df_transformed.show()

Вывод:

+-----------+-------+-----+----+----------+----------+----------------+-------------+
|   category|product|sales|rank|dense_rank|row_number|cumulative_sales|average_sales|
+-----------+-------+-----+----+----------+----------+----------------+-------------+
|Electronics| Tablet|  800|   1|         1|         1|             800|        800.0|
|Electronics|  Phone| 1000|   2|         2|         2|            1800|        900.0|
|Electronics| Laptop| 1500|   3|         3|         3|            3300|       1100.0|
|  Furniture|  Chair|  300|   1|         1|         1|             600|        300.0|
|  Furniture|  Table|  300|   1|         1|         2|             600|        300.0|
|  Furniture|   Desk|  600|   3|         2|         3|            1200|        400.0|
+-----------+-------+-----+----+----------+----------+----------------+-------------+

Итог

В этом уроке рассмотрены:

  • Оконные спецификации с partitionBy и orderBy
  • Агрегатные функции (sum, avg, min, max, count) для нарастающих итогов и средних значений
  • Ранжирующие функции (rank, dense_rank, row_number) для присвоения рангов и порядковых номеров