Aller au contenu principal

agréger

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 finalisation. Prend en charge Spark Connect.

Pour la fonction Databricks SQL correspondante, consultez la fonctionaggregate.

Syntaxe

Python
from pyspark.sql import functions as dbf

dbf.aggregate(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 retourne une expression du même type que initialValue.

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 retourne une expression du même type que initialValue.

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** : Agrégation simple 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.aggregate("values", dbf.lit(0.0), lambda acc, x: acc + x).alias("sum")).show()
Output
+----+
| sum|
+----+
|42.0|
+----+

**Exemple 2** : Agrégation 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.aggregate(
"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|
+----+