réduire
Applique un opérateur binaire à un état initial et à tous les éléments du tableau, et réduit cela à un seul état. L'état final est converti en résultat final en appliquant une fonction de finition.
Pour la fonction Databricks SQL correspondante, consultez la fonctionreduce.
Syntaxe
from pyspark.sql import functions as dbf
dbf.reduce(col=<col>, initialValue=<initialValue>, merge=<merge>, finish=<finish>)
parameter
parameter | Type | Description |
|---|---|---|
|
| Nom de la colonne ou de l'expression. |
|
| Valeur initiale. Nom de la colonne ou expression. |
|
| Une fonction binaire qui renvoie une expression du même type que zéro. |
|
| Une fonction unaire facultative utilisée pour convertir une valeur accumulée. |
Renvoie
pyspark.sql.Column: valeur finale après l'application de la fonction d'agrégation.
Exemples
Exemple 1 : simple réduction avec somme
from pyspark.sql import functions as dbf
df = spark.createDataFrame([(1, [20.0, 4.0, 2.0, 6.0, 10.0])], ("id", "values"))
df.select(dbf.reduce("values", dbf.lit(0.0), lambda acc, x: acc + x).alias("sum")).show()
+----+
| sum|
+----+
|42.0|
+----+
Exemple 2 : Réduction avec fonction de fin
from pyspark.sql import functions as dbf
df = spark.createDataFrame([(1, [20.0, 4.0, 2.0, 6.0, 10.0])], ("id", "values"))
def merge(acc, x):
count = acc.count + 1
sum = acc.sum + x
return dbf.struct(count.alias("count"), sum.alias("sum"))
df.select(
dbf.reduce(
"values",
dbf.struct(dbf.lit(0).alias("count"), dbf.lit(0.0).alias("sum")),
merge,
lambda acc: acc.sum / acc.count,
).alias("mean")
).show()
+----+
|mean|
+----+
| 8.4|
+----+