如何在PySpark中无需聚合,按ID分组实现手机号列的表透视?
Spark实现按ID分组透视手机号为多列的方案
可以通过**窗口函数+透视(pivot)**实现需求,无需真正的聚合操作(仅语法上需要聚合函数,实际是取唯一值),具体步骤如下:
步骤说明
- 添加分组内序号:用窗口函数给每个ID下的手机号分配唯一序号(1、2、3),确保每个手机号对应一个序号;
- 透视生成独立列:以序号为透视列,将不同序号的手机号转为独立列,用
first()取对应值(因每个(ID,序号)组合仅一条数据,无聚合逻辑)。
代码实现
Python版本
from pyspark.sql import Window import pyspark.sql.functions as F # 假设原始数据表对应的DataFrame为df,结构为ID(string)、Phone(string) # 定义窗口规则:按ID分区,手机号排序(可根据需求调整排序字段) window_spec = Window.partitionBy("ID").orderBy("Phone") # 给每个ID下的手机号添加分组内序号 df_with_rn = df.withColumn("rn", F.row_number().over(window_spec)) # 透视生成目标列,指定序号范围为1-3,避免自动生成多余列 result_df = df_with_rn.groupBy("ID").pivot("rn", [1, 2, 3]).agg(F.first("Phone")) # 重命名列名,匹配预期结果的Phone1/Phone2/Phone3 result_df = result_df.withColumnRenamed("1", "Phone1") \ .withColumnRenamed("2", "Phone2") \ .withColumnRenamed("3", "Phone3") # 查看结果 result_df.show()
Scala版本
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // 假设原始数据表对应的DataFrame为df,结构为ID(String)、Phone(String) val windowSpec = Window.partitionBy("ID").orderBy("Phone") // 添加分组内序号 val dfWithRn = df.withColumn("rn", row_number().over(windowSpec)) // 透视并生成目标列,指定序号范围1-3 val resultDf = dfWithRn.groupBy("ID") .pivot("rn", Seq(1, 2, 3)) .agg(first("Phone")) .withColumnRenamed("1", "Phone1") .withColumnRenamed("2", "Phone2") .withColumnRenamed("3", "Phone3") // 查看结果 resultDf.show()
关键细节
- 排序规则:
orderBy("Phone")可根据实际需求调整(比如按数据插入顺序排序,可改用monotonically_increasing_id()),只要保证每个ID下的手机号序号唯一即可; - 空值处理:若某个ID的手机号少于3个,对应列会显示
null,与预期结果一致; - 聚合函数:这里用
first()是因为每个(ID,序号)组合仅存在一条数据,不存在聚合计算,用max()/min()也能得到相同结果。
内容的提问来源于stack exchange,提问作者Guilherme Mendes
相关产品推荐
相关产品推荐

