aggregate

İlk duruma ve dizideki tüm öğelere ikili işleç uygular ve bunu tek bir duruma küçültür. Son durum, bir bitiş işlevi uygulanarak son sonucta dönüştürülür. Spark Connect'i destekler.

İlgili Databricks SQL fonksiyonu için bakınız aggregate fonksiyonu.

Sözdizimi

from pyspark.sql import functions as dbf

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

Parametreler

Parametre Türü Description
col pyspark.sql.Column veya str Sütun veya ifadenin adı.
initialValue pyspark.sql.Column veya str İlk değer. Sütun veya ifadenin adı.
merge function initialValue ile aynı türde bir ifade döndüren ikili işlev.
finish functionopsiyonel Birikmiş değeri dönüştürmek için kullanılan isteğe bağlı bir tekli işlev.

İade

pyspark.sql.Column: toplama işlevi uygulandıktan sonraki son değer.

Örnekler

Örnek 1: Toplam ile basit toplama

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|
+----+

Örnek 2: Bitiş işleviyle toplama

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|
+----+