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

PySpark分组透视后将多列转为Map的实现问题及代码修正

PySpark分组透视与Map类型转换问题修复

需求说明

对DataFrame按ID分组,基于Order和TYPE字段生成透视列,同时将Address、Phone列合并为Map类型字段。

输入DataFrame

IDTYPEOrderAddressPhone
1APrimaryabc111
1BPrimarydef222
1ASecondaryghi333
1BSecondaryjkl444
2APrimarymno555
2ASecondarypqr666
2BPrimarystu777
2BSecondaryvwy888

预期输出DataFrame

IDPrimary_A_attributesPrimary_B_attributesSecondary_A_attributesSecondary_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()

问题分析

  1. 列名拼接错误:原代码用_attributes_作为Order和TYPE的分隔符,生成的列名格式为PRIMARY_ATTRIBUTES_A,与预期的Primary_A_attributes不符。
  2. 不必要的大写转换:upper()将列名转为全大写,不符合预期的首字母大写格式。
  3. 聚合逻辑潜在风险:未指定明确的聚合函数,虽当前数据每个分组唯一,但显式指定更稳妥。

修复后的代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 01:18:27