如何在PySpark中按行分组并创建新列(适配大文件)
问题描述
原始DataFrame
| id | name | |
|---|---|---|
| 1 | id1@first.com | john |
| 2 | id2@first.com | Maike |
| 2 | id2@second | Maike |
| 1 | id1@second.com | john |
目标转换结果
| id | email1 | name | |
|---|---|---|---|
| 1 | id1@first.com | id1@second.com | john |
| 2 | id2@first.com | id2@second | Maike |
实际待处理的文件包含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
相关产品推荐
相关产品推荐

