計算 聚合,並回傳結果為 DataFrame。
可用的聚合函數可以有:
- 內建的聚合函數,如
avg、max、minsumcount、 , 。 - 群組聚合 pandas UDF,使用
pyspark.sql.functions.pandas_udf。
語法
agg(*exprs)
參數
| 參數 | 類型 | 說明 |
|---|---|---|
exprs |
字條或專欄 | 從欄位名稱(字串)映射到聚合函數(字串)的字典,或是聚合 Column 表達式的清單。 |
退貨
DataFrame
Notes
內建的聚合函數與群組聚合 PANDAS UDF 不能在同一次呼叫中混合使用。
當 exprs 是單一字令時,鍵為要執行彙總的欄位,值為彙總函數。 當 exprs 為表達式列表 Column 時,每個表達式指定一個要計算的聚合。
Examples
import pandas as pd
from pyspark.sql import functions as sf
df = spark.createDataFrame(
[(2, "Alice"), (3, "Alice"), (5, "Bob"), (10, "Bob")], ["age", "name"])
# Group-by name, and count each group.
df.groupBy(df.name).agg({"*": "count"}).sort("name").show()
# +-----+--------+
# | name|count(1)|
# +-----+--------+
# |Alice| 2|
# | Bob| 2|
# +-----+--------+
# Group-by name, and calculate the minimum age.
df.groupBy(df.name).agg(sf.min(df.age)).sort("name").show()
# +-----+--------+
# | name|min(age)|
# +-----+--------+
# |Alice| 2|
# | Bob| 5|
# +-----+--------+
# Same as above but uses a pandas UDF.
from pyspark.sql.functions import pandas_udf
@pandas_udf('int')
def min_udf(v: pd.Series) -> int:
return v.min()
df.groupBy(df.name).agg(min_udf(df.age)).sort("name").show()
# +-----+------------+
# | name|min_udf(age)|
# +-----+------------+
# |Alice| 2|
# | Bob| 5|
# +-----+------------+