Aller au contenu principal

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

Python
from pyspark.sql import functions as dbf

dbf.reduce(col=<col>, initialValue=<initialValue>, merge=<merge>, finish=<finish>)

parameter

parameter

Type

Description

col

pyspark.sql.Column OU str

Nom de la colonne ou de l'expression.

initialValue

pyspark.sql.Column OU str

Valeur initiale. Nom de la colonne ou expression.

merge

function

Une fonction binaire qui renvoie une expression du même type que zéro.

finish

function, facultatif

Une fonction unaire facultative utilisée pour convertir une valeur accumulée.

parameter

Type

Description

col

pyspark.sql.Column OU str

Nom de la colonne ou de l'expression.

initialValue

pyspark.sql.Column OU str

Valeur initiale. Nom de la colonne ou expression.

merge

function

Une fonction binaire qui renvoie une expression du même type que zéro.

finish

function, facultatif

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

Python
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()
Output
+----+
| sum|
+----+
|42.0|
+----+

Exemple 2 : Réduction avec fonction de fin

Python
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()
Output
+----+
|mean|
+----+
| 8.4|
+----+