Класс Window

Служебные функции для определения окна в кадрах данных.

Поддержка Spark Connect

Атрибуты классов

Атрибут Описание
unboundedPreceding Значение границы, представляющее начало необвязанной рамки окна.
unboundedFollowing Значение границы, представляющее конец несвязанной рамки окна.
currentRow Значение границы, представляющее текущую строку в рамке окна.

Методы

Метод Описание
orderBy(*cols) Создает WindowSpec с определенным упорядочением.
partitionBy(*cols) Создает WindowSpec с определенным секционированием.
rangeBetween(start, end) Создает WindowSpec с заданными границами кадра от start (включительно) до end (включительно), используя смещения на основе диапазона из значения текущей строки ORDER BY .
rowsBetween(start, end) Создает WindowSpec с заданными границами фрейма (от start (включительно) до end (включительно), используя смещения на основе строк из текущей строки.

Примечания

Если порядок не определен, по умолчанию используется необвязанная рамка окна (rowFrame, unboundedFollowing, unboundedFollowing). При определении упорядочения по умолчанию используется растущий кадр окна (rangeFrame, unbounded, currentRow).

Примеры

Базовое окно с упорядочением и кадром строк

from pyspark.sql import Window

# ORDER BY date ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
window = Window.orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow)

Секционированное окно с кадром диапазона

from pyspark.sql import Window

# PARTITION BY country ORDER BY date RANGE BETWEEN 3 PRECEDING AND 3 FOLLOWING
window = Window.orderBy("date").partitionBy("country").rangeBetween(-3, 3)

Номер строки в секции

from pyspark.sql import Window, functions as sf

df = spark.createDataFrame(
    [(1, "a"), (1, "a"), (2, "a"), (1, "b"), (2, "b"), (3, "b")], ["id", "category"]
)

# Show row number ordered by id within each category partition
window = Window.partitionBy("category").orderBy("id")
df.withColumn("row_number", sf.row_number().over(window)).show()

Выполнение суммы с кадром на основе строк

from pyspark.sql import Window, functions as sf

df = spark.createDataFrame(
    [(1, "a"), (1, "a"), (2, "a"), (1, "b"), (2, "b"), (3, "b")], ["id", "category"]
)

# Sum id values from the current row to the next row within each partition
window = Window.partitionBy("category").orderBy("id").rowsBetween(Window.currentRow, 1)
df.withColumn("sum", sf.sum("id").over(window)).sort("id", "category", "sum").show()

Выполнение суммы с кадром на основе диапазона

from pyspark.sql import Window, functions as sf

df = spark.createDataFrame(
    [(1, "a"), (1, "a"), (2, "a"), (1, "b"), (2, "b"), (3, "b")], ["id", "category"]
)

# Sum id values from the current id value to id + 1 within each partition
window = Window.partitionBy("category").orderBy("id").rangeBetween(Window.currentRow, 1)
df.withColumn("sum", sf.sum("id").over(window)).sort("id", "category").show()