如何在PySpark DataFrame中对名称相似的列计算平均值
如何在PySpark DataFrame中对名称相似的列计算平均值
嗨~看你的需求是要给PySpark DataFrame里的12个月余额列(b0到b11)新增一个平均余额列,而且想要通用的代码,不用硬写所有列名对吧?这就给你安排两个实用的方法,都能灵活适配列名规律的情况~
方法一:用内置数组平均值函数(简洁通用)
这个方法利用PySpark的内置函数,自动识别符合命名规律的列,并且会自动忽略空值(null),完全匹配你例子里的计算逻辑。
步骤:
- 先导入PySpark的函数库
- 筛选出所有目标列(比如符合
b+数字的命名规则) - 用数组平均值函数计算每行的平均余额
from pyspark.sql import functions as F import re # 假设你的DataFrame名为df # 精准筛选列名格式为b+数字的列(比如b0、b1...b11) balance_cols = [col for col in df.columns if re.match(r'^b\d+$', col)] # 新增avg_bal列,计算每行目标列的平均值(自动忽略null) df_with_avg = df.withColumn( 'avg_bal', F.expr(f"avg(array({','.join(balance_cols)}))") ) # 查看结果 df_with_avg.show()
说明:
re.match(r'^b\d+$', col)用来精准匹配列名,确保只选中b0到b11这类列,避免误选其他含b的列avg(array(...))会对每行的数组元素求平均值,自动跳过null值,和你例子里的计算结果完全一致:比如cust_1的(20+30)/2=25,cust_3的(50+30+10)/3=30,全null的cust_4会返回null
方法二:手动聚合计算(自定义逻辑更灵活)
如果你需要更灵活的处理逻辑(比如自定义空值的处理方式、添加权重等),可以手动计算总和和非空值的数量,再求平均。
from pyspark.sql import functions as F import re balance_cols = [col for col in df.columns if re.match(r'^b\d+$', col)] df_with_avg = df.withColumn( 'avg_bal', F.aggregate( # 将目标列转换成数组 F.array(*balance_cols), # 初始化累加器:sum存储总和,count存储非空值的数量 F.struct(F.lit(0).alias('sum'), F.lit(0).alias('count')), # 遍历数组元素,更新总和和计数 lambda acc, x: F.struct( acc.sum + F.coalesce(x, F.lit(0)), acc.count + F.when(F.isnull(x), F.lit(0)).otherwise(F.lit(1)) ), # 最后计算平均值:如果计数为0(全null)返回null,否则返回总和/计数 lambda acc: F.when(acc.count == 0, F.lit(None)).otherwise(acc.sum / acc.count) ) ) df_with_avg.show()
说明:
- 这个方法完全自定义计算逻辑,你可以根据需求修改累加规则,比如给某些月份的余额加权重,或者把全null的情况返回0而不是
null - 同样能自动适配所有符合命名规则的列,不用手动列全12个月份的列名
总结
如果只是简单计算忽略空值的行级平均值,方法一足够简洁通用;如果需要特殊的自定义逻辑,方法二更适合。两种方法都不需要硬编码所有列名,完全符合你要的“通用代码”需求~
备注:内容来源于stack exchange,提问作者suraj jadhav
相关产品推荐
相关产品推荐

