You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.23 21:36:04