【发布时间】:2020-11-26 13:52:14
【问题描述】:
我对 Pyspark 中的 udfs 和一个具体案例有疑问。 我正在尝试制作一个简单的可重用函数来聚合不同级别和组的值。 输入应该是:
- 现有数据框
- 分组依据的变量(单列或列表)
- 要聚合的变量(同上)
- 要应用的函数(一个特定的函数或它们的列表)。我保持简单的求和、平均、最小值、最大值等......
当我有一个函数或一个列表时,我让它工作,但是当涉及到聚合变量时,我被困在将它们的列表引入函数中
def aggregate(dataframe,grouping,aggregation,functions):
**First part works ok on single functions and single columns**
if hasattr(aggregation,'__iter__') == False and hasattr(functions,'__iter__') == False:
if functions == sum:
df = dataframe.groupby(grouping).sum(aggregation)
elif functions == avg:
df = dataframe.groupby(grouping).avg(aggregation)
elif functions == min:
df = dataframe.groupby(grouping).min(aggregation)
elif functions == max:
df = dataframe.groupby(grouping).max(aggregation)
elif functions == count:
df = dataframe.groupby(grouping).count(aggregation)
elif functions == countDistinct:
df = dataframe.groupby(grouping).countDistinct(aggregation)
**Here is where I got into the part I struggle with, if aggregation == [some list] it will not work
elif hasattr(aggregation,'__iter__') == True and hasattr(functions,'__iter__') == False:
if functions == sum:
df = dataframe.groupby(grouping).sum(aggregation)
elif functions == avg:
df = dataframe.groupby(grouping).avg(aggregation)
elif functions == min:
df = dataframe.groupby(grouping).min(aggregation)
elif functions == max:
df = dataframe.groupby(grouping).max(aggregation)
elif functions == count:
df = dataframe.groupby(grouping).count(aggregation)
elif functions == countDistinct:
df = dataframe.groupby(grouping).countDistinct(aggregation)
**Expression to get inputs as lists works too**
else:
expression_def = [f(col(c)) for f in functions for c in aggregation]
df = dataframe.groupby(grouping).agg(*expression_def)
return df
【问题讨论】:
标签: python apache-spark pyspark apache-spark-sql