如何在PySpark中对同类型列进行批量求和运算?
PySpark 批量同前缀列求和实现
以下是通用的代码实现,能自动识别同前缀列(如b和b_apac)并完成求和,最终保留id和所有求和后的新列:
from pyspark.sql import functions as F # 替换为你的实际DataFrame # df = spark.read.table("your_table") # 示例DataFrame(仅作演示) data = [ (1, 2, 3, 4, 5, 6), (2, 7, 8, 9, 10, 11) ] columns = ["id", "b", "b_apac", "c", "c_apac", "d", "d_apac"] df = spark.createDataFrame(data, schema=columns) # 1. 按前缀分组非id列 non_id_cols = [col for col in df.columns if col != "id"] prefix_groups = {} for col in non_id_cols: # 提取前缀:处理基础列和_apac后缀列 prefix = col.split("_apac")[0] if "_apac" in col else col if prefix not in prefix_groups: prefix_groups[prefix] = [] prefix_groups[prefix].append(col) # 2. 生成求和表达式 sum_columns = [] for prefix, cols in prefix_groups.items(): # 对同前缀列求和,命名为「前缀_sum」 # 两种求和方式可选: # 方式1:适合任意数量的同前缀列(比如后续加b_emea也能自动求和) sum_expr = F.sum(F.array(*cols)).alias(f"{prefix}_sum") # 方式2:仅适用于固定两列的场景(基础列+_apac列) # sum_expr = (F.col(cols[0]) + F.col(cols[1])).alias(f"{prefix}_sum") sum_columns.append(sum_expr) # 3. 生成结果表:保留id + 所有求和列 result_df = df.select("id", *sum_columns) # 查看输出 result_df.show()
关键说明
- 自动分组:通过识别
_apac后缀提取前缀,不管是基础列(如b)还是扩展列(如b_apac),都会被归到同一前缀组,无需手动枚举所有列 - 灵活求和:提供两种求和方式,若后续可能新增其他区域列(如
b_emea),推荐用F.sum(F.array(*cols))的方式,能自动适配多列求和 - 通用性:不管有多少列(哪怕超过100列),只要符合命名规则,代码都能自动处理,无需修改核心逻辑
内容的提问来源于stack exchange,提问作者suraj jadhav
相关产品推荐
相关产品推荐

