如何使用PySpark通过groupby对分组后多列行值进行求和
PySpark 分组多列求和实现方案
需求说明
现有PySpark DataFrame需按以下规则完成计算:
- 分组字段:
load_dt、org_cntry - 聚合逻辑:每个分组内对
sum(srv_curr_vo_qty_accs_mthd)、sum(srv_curr_bb_qty_accs_mthd)、sum(srv_curr_tv_qty_accs_mthd)三个列的数值求和,生成新字段total_sum
原始数据示例
load_dt|org_cntry|sum(srv_curr_vo_qty_accs_mthd)|sum(srv_curr_bb_qty_accs_mthd)|sum(srv_curr_tv_qty_accs_mthd)| +-------------------+---------+------------------------------+------------------------------+------------------------------+ |2021-12-06 00:00:00| null| NaN| NaN| NaN| |2021-12-06 00:00:00| PANAMA| 360126.0| 214229.0| 207950.0|
实现代码
from pyspark.sql import functions as F # 替换下方df为你实际的原始DataFrame变量名 result_df = df \ # 过滤org_cntry为空的无效行 .filter(F.col("org_cntry").isNotNull()) \ # 格式调整:转换load_dt为日期类型,org_cntry统一为首字母大写格式 .withColumn("load_dt", F.to_date(F.col("load_dt"))) \ .withColumn("org_cntry", F.initcap(F.lower(F.col("org_cntry")))) \ # 按指定字段分组 .groupBy("load_dt", "org_cntry") \ # 三列求和生成结果字段 .agg( (F.sum("sum(srv_curr_vo_qty_accs_mthd)") + F.sum("sum(srv_curr_bb_qty_accs_mthd)") + F.sum("sum(srv_curr_tv_qty_accs_mthd)") ).cast("int").alias("total_sum") ) # 输出验证结果 result_df.show()
输出结果示例
+----------+---------+---------+ | load_dt|org_cntry|total_sum| +----------+---------+---------+ |2021-12-06| Panama| 782305| +----------+---------+---------+
内容的提问来源于stack exchange,提问作者Pavithra Kannan
相关产品推荐
相关产品推荐

