Функции 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 — требуют оконной спецификации, определяющей:
- Партиционирование: как набор данных разбивается на группы
- Упорядочивание: как строки внутри каждой партиции упорядочиваются
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|
+------------+----------+------+-----------+--------------+------------+