délai
Fonction de fenêtre : renvoie la valeur qui se trouve à offset lignes avant la ligne actuelle, et default s'il y a moins de offset lignes avant la ligne actuelle. Par exemple, un offset de un renverra la ligne précédente à tout moment donné dans la partition de fenêtre.
C'est l'équivalent de la fonction LAG en SQL.
Syntaxe
from pyspark.sql import functions as sf
sf.lag(col, offset=1, default=None)
parameter
parameter | Type | Description |
|---|---|---|
|
| Nom de la colonne ou de l'expression. |
| int, facultatif | Nombre de lignes à étendre. La valeur par default est 1. |
| Facultatif | Valeur default. |
Renvoie
pyspark.sql.Column: valeur avant la ligne actuelle basée sur offset.
Exemples
Exemple 1 : Utilisation de lag pour obtenir la valeur précédente
from pyspark.sql import functions as sf
from pyspark.sql import Window
df = spark.createDataFrame(
[("a", 1), ("a", 2), ("a", 3), ("b", 8), ("b", 2)], ["c1", "c2"])
df.show()
+---+---+
| c1| c2|
+---+---+
| a| 1|
| a| 2|
| a| 3|
| b| 8|
| b| 2|
+---+---+
w = Window.partitionBy("c1").orderBy("c2")
df.withColumn("previous_value", sf.lag("c2").over(w)).show()
+---+---+--------------+
| c1| c2|previous_value|
+---+---+--------------+
| a| 1| NULL|
| a| 2| 1|
| a| 3| 2|
| b| 2| NULL|
| b| 8| 2|
+---+---+--------------+
**Exemple 2** : Utilisation du décalage avec une valeur par default
from pyspark.sql import functions as sf
from pyspark.sql import Window
df = spark.createDataFrame(
[("a", 1), ("a", 2), ("a", 3), ("b", 8), ("b", 2)], ["c1", "c2"])
w = Window.partitionBy("c1").orderBy("c2")
df.withColumn("previous_value", sf.lag("c2", 1, 0).over(w)).show()
+---+---+--------------+
| c1| c2|previous_value|
+---+---+--------------+
| a| 1| 0|
| a| 2| 1|
| a| 3| 2|
| b| 2| 0|
| b| 8| 2|
+---+---+--------------+
**Exemple 3** : Utilisation d'un décalage avec un offset de 2
from pyspark.sql import functions as sf
from pyspark.sql import Window
df = spark.createDataFrame(
[("a", 1), ("a", 2), ("a", 3), ("b", 8), ("b", 2)], ["c1", "c2"])
w = Window.partitionBy("c1").orderBy("c2")
df.withColumn("previous_value", sf.lag("c2", 2, -1).over(w)).show()
+---+---+--------------+
| c1| c2|previous_value|
+---+---+--------------+
| a| 1| -1|
| a| 2| -1|
| a| 3| 1|
| b| 2| -1|
| b| 8| -1|
+---+---+--------------+