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

如何在PySpark中按行分组并创建新列(适配大文件)

问题描述

原始DataFrame

idemailname
1id1@first.comjohn
2id2@first.comMaike
2id2@secondMaike
1id1@second.comjohn

目标转换结果

idemailemail1name
1id1@first.comid1@second.comjohn
2id2@first.comid2@secondMaike

实际待处理的文件包含60余列,数据量庞大。目前使用spark.read读取文件(优先选择该方式,因为支持延迟计算特性,而pyspark.pandas API不支持),需要实现按id分组,将同一分组下的多个email拆分为单独列,同时保留其他列的转换逻辑。

解决方案

1. 读取数据

修正文件名拼写后,使用spark.read读取CSV文件:

df = spark.read.option("header", True) \
        .csv("contacts.csv", sep=',')

2. 为分组内的email添加序号

使用窗口函数row_number(),按id分组后给每个email分配序号,方便后续行转列:

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, col

# 定义窗口:按id分组,可根据需求指定排序规则(比如按email字母顺序)
window_spec = Window.partitionBy("id").orderBy("email")

# 添加序号列
df_with_rank = df.withColumn("email_rank", row_number().over(window_spec))

3. 行转列实现分组聚合

利用pivot将同一id下的多个email转换成独立列,同时聚合其他列(同一id下值一致的列直接取首个值即可):

# 分组后进行pivot,指定email序号范围(若不确定数量可省略[1,2],但指定范围能提升性能)
pivoted_df = df_with_rank.groupBy("id", "name") \
    .pivot("email_rank", [1, 2]) \
    .agg({"email": "first"}) \
    .withColumnRenamed("1", "email") \
    .withColumnRenamed("2", "email1")

4. 调整列顺序(可选)

若需要和目标结构列顺序完全一致,可手动指定列顺序:

target_columns = ["id", "email", "email1", "name"]
final_df = pivoted_df.select(target_columns)

# 查看转换结果
final_df.show()
注意事项
  • 若同一id下的email数量超过2个,需调整pivot中的序号列表(比如[1,2,3]),或不指定列表让Spark自动识别;数据量较大时,建议提前统计分组内的最大email数量,指定序号能大幅提升执行性能
  • 对于其他60余列,只要同一id下值一致,均可加入groupBy列表,聚合时使用first()/max()等函数;若存在同一id下值不同的列,需根据业务逻辑确定聚合方式
  • 所有转换操作均遵循Spark的延迟计算特性,只有触发show()、write()等动作时才会实际执行计算

内容的提问来源于stack exchange,提问作者Emerson Pedroso

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 01:15:35