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
from pyspark.sql import functions as dbf
dbf.aggregate(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 retourne une expression du même type que initialValue. |
|
| 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
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()
+----+
| sum|
+----+
|42.0|
+----+
**Exemple 2** : Agrégation 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.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()
+----+
|mean|
+----+
| 8.4|
+----+