PySpark分组透视后将多列转为Map的实现问题及代码修正
PySpark分组透视与Map类型转换问题修复
需求说明
对DataFrame按ID分组,基于Order和TYPE字段生成透视列,同时将Address、Phone列合并为Map类型字段。
输入DataFrame
| ID | TYPE | Order | Address | Phone |
|---|---|---|---|---|
| 1 | A | Primary | abc | 111 |
| 1 | B | Primary | def | 222 |
| 1 | A | Secondary | ghi | 333 |
| 1 | B | Secondary | jkl | 444 |
| 2 | A | Primary | mno | 555 |
| 2 | A | Secondary | pqr | 666 |
| 2 | B | Primary | stu | 777 |
| 2 | B | Secondary | vwy | 888 |
预期输出DataFrame
| ID | Primary_A_attributes | Primary_B_attributes | Secondary_A_attributes | Secondary_B_attributes |
|---|---|---|---|---|
| 1 | {"Address":"abc","phone":"111"} | {"Address":"def","phone":"222"} | {"Address":"ghi","phone":"333"} | {"Address":"jkl","phone":"444"} |
| 2 | {"Address":"mno","phone":"555"} | {"Address":"pqr","phone":"666"} | {"Address":"stu","phone":"777"} | {"Address":"vwy","phone":"888"} |
问题代码
df.withColumn("collection",F.upper(F.concat_ws('_attributes_','order','type')))\ .groupBy('id').pivot("collection").agg(F.create_map(F.lit("Address"),F.col("Address"),\ F.lit("phone"),F.col("phone"))).display()
问题分析
- 列名拼接错误:原代码用
_attributes_作为Order和TYPE的分隔符,生成的列名格式为PRIMARY_ATTRIBUTES_A,与预期的Primary_A_attributes不符。 - 不必要的大写转换:
upper()将列名转为全大写,不符合预期的首字母大写格式。 - 聚合逻辑潜在风险:未指定明确的聚合函数,虽当前数据每个分组唯一,但显式指定更稳妥。
修复后的代码
import pyspark.sql.functions as F df.withColumn( "collection", F.concat(F.col("Order"), F.lit("_"), F.col("TYPE"), F.lit("_attributes")) ).groupBy("ID").pivot("collection").agg( F.first(F.create_map( F.lit("Address"), F.col("Address"), F.lit("phone"), F.col("Phone") )) ).display()
修复说明
- 修正列名生成逻辑:通过
concat()拼接Order、TYPE和固定后缀_attributes,生成符合预期的列名格式。 - 移除大写转换:保留原字段的首字母大写格式,与预期输出一致。
- 添加明确聚合函数:使用
first()确保每个透视列取到对应分组的唯一Map值,避免潜在的多行聚合问题。
内容的提问来源于stack exchange,提问作者pradeep nadarajan
相关产品推荐
相关产品推荐

