Функции lead() и lag() в PySpark

lead() для доступа к следующей строке, lag() для доступа к предыдущей строке в окне: параметры column/offset/default, вычисление разниц, предсказание событий, обработка граничных значений.

core

Понимание lead и lag

lead(column, offset, default)

Получает значение указанной колонки из следующей строки в пределах одного окна.

Параметры:

  • column: Колонка, из которой берётся значение
  • offset: Количество строк для просмотра вперёд (по умолчанию 1)
  • default: Значение, возвращаемое при выходе за пределы окна

Применение:

  • Сравнение текущего и следующего значений для выявления трендов роста или спада
  • Предсказание следующей транзакции или события для пользователя или объекта

lag(column, offset, default)

Получает значение указанной колонки из предыдущей строки в пределах одного окна.

Параметры:

  • column: Колонка, из которой берётся значение
  • offset: Количество строк для просмотра назад (по умолчанию 1)
  • default: Значение, возвращаемое при выходе за пределы окна

Применение:

  • Сравнение текущего и предыдущего значений для вычисления разниц и трендов
  • Получение исторического контекста для объекта или события

Оконная спецификация

Обе функции — lead и lag — требуют оконной спецификации, определяющей:

  1. Партиционирование: как набор данных разбивается на группы
  2. Упорядочивание: как строки внутри каждой партиции упорядочиваются
from pyspark.sql.window import Window
window_spec = Window.partitionBy("group_column").orderBy("ordering_column")

Практические примеры

Пример 1: Следующая и предыдущая суммы транзакции

from pyspark.sql import functions as F
from pyspark.sql.window import Window

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

# Добавление колонок lead и lag
df = df.withColumn("next_transaction", F.lead("amount", 1).over(window_spec))
df = df.withColumn("previous_transaction", F.lag("amount", 1).over(window_spec))

Вывод:

+------------+----------+------+-----------------+---------------------+
| customer_id|      date|amount|next_transaction |previous_transaction |
+------------+----------+------+-----------------+---------------------+
|           1|2024-01-01|   100|              200|                 null|
|           1|2024-01-02|   200|              300|                  100|
|           1|2024-01-03|   300|             null|                  200|
|           2|2024-01-01|    50|               80|                 null|
|           2|2024-01-02|    80|              120|                   50|
|           2|2024-01-03|   120|             null|                   80|
+------------+----------+------+-----------------+---------------------+

Пример 2: Вычисление разницы между покупками

# Разница между текущей и предыдущей транзакцией
df = df.withColumn("purchase_difference", F.col("amount") - F.lag("amount", 1).over(window_spec))

Вывод:

+------------+----------+------+-------------------+
| customer_id|      date|amount|purchase_difference|
+------------+----------+------+-------------------+
|           1|2024-01-01|   100|               null|
|           1|2024-01-02|   200|                100|
|           1|2024-01-03|   300|                100|
|           2|2024-01-01|    50|               null|
|           2|2024-01-02|    80|                 30|
|           2|2024-01-03|   120|                 40|
+------------+----------+------+-------------------+

Пример 3: Предсказание следующего события

df = df.withColumn("next_status", F.lead("status", 1).over(window_spec))

Вывод:

+------------+----------+----------+--------------+
| customer_id|      date|    status|   next_status|
+------------+----------+----------+--------------+
|           1|2024-01-01|   browsing|       carted|
|           1|2024-01-02|     carted|    purchased|
|           1|2024-01-03|  purchased|         null|
+------------+----------+----------+--------------+

Пример 4: Подстановка значения по умолчанию вместо null

df = df.withColumn("previous_transaction", F.lag("amount", 1, 0).over(window_spec))

Вывод:

+------------+----------+------+--------------------+
| customer_id|      date|amount|previous_transaction|
+------------+----------+------+--------------------+
|           1|2024-01-01|   100|                   0|
|           1|2024-01-02|   200|                 100|
|           1|2024-01-03|   300|                 200|
+------------+----------+------+--------------------+

Пример 5: Обнаружение изменений статуса

df = df.withColumn("status_change", F.col("status") != F.lag("status", 1).over(window_spec))

Вывод:

+------------+----------+----------+-------------+
| customer_id|      date|    status|status_change|
+------------+----------+----------+-------------+
|           1|2024-01-01| browsing |        false|
|           1|2024-01-02|   carted |         true|
|           1|2024-01-03| purchased|         true|
+------------+----------+----------+-------------+

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

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window

# Тестовые данные
data = [
    (1, "2024-01-01", 100),
    (1, "2024-01-02", 200),
    (1, "2024-01-03", 300),
    (2, "2024-01-01", 50),
    (2, "2024-01-02", 80),
    (2, "2024-01-03", 120),
]

# Создание DataFrame
spark = SparkSession.builder.appName("LeadLagTutorial").getOrCreate()
df = spark.createDataFrame(data, ["customer_id", "date", "amount"])

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

# Применение lead и lag
df = df.withColumn("next_amount", F.lead("amount", 1).over(window_spec))
df = df.withColumn("previous_amount", F.lag("amount", 1).over(window_spec))
df = df.withColumn("amount_diff", F.col("amount") - F.lag("amount", 1).over(window_spec))

df.show()

Вывод:

+------------+----------+------+-----------+--------------+------------+
| customer_id|      date|amount|next_amount|previous_amount|amount_diff|
+------------+----------+------+-----------+--------------+------------+
|           1|2024-01-01|   100|        200|          null|        null|
|           1|2024-01-02|   200|        300|           100|         100|
|           1|2024-01-03|   300|       null|           200|         100|
|           2|2024-01-01|    50|         80|          null|        null|
|           2|2024-01-02|    80|        120|            50|          30|
|           2|2024-01-03|   120|       null|            80|          40|
+------------+----------+------+-----------+--------------+------------+