在Databricks PySpark中转置DataFrame并合并CxName列
Databricks PySpark DataFrame转置并拼接CxName列
核心实现步骤
要达成转置后CxName列为逗号分隔的拼接字符串,需先完成分组拼接,再执行转置,具体操作如下:
1. 分组拼接CxName
先按Acct分组,将同一账户下的CxName按Position排序后拼接成字符串:
from pyspark.sql import functions as F # 原始数据 df = spark.createDataFrame([('12345','John Doe',1),('12345','Jane Doe',2)],['Acct' , 'CxName' ,'Position' ]) # 按Acct分组,按Position排序后拼接CxName df_concat = df.groupBy("Acct") \ .agg(F.concat_ws(", ", F.collect_list(F.struct("Position", "CxName")).orderBy("Position").getField("CxName")).alias("CxName"))
2. 执行转置操作
提供两种转置实现方式,按需选择:
方式一:基于pyspark.pandas转置
import pyspark.pandas as ps # 转置并调整列名 transposed_df = df_concat.pandas_api() \ .set_index("Acct") \ .T \ .reset_index() \ .rename(columns={"index": "Column"}) \ .to_spark() # 查看结果 transposed_df.show()
方式二:纯PySpark原生函数转置(无pandas依赖)
# 创建键值对映射并展开实现转置 transposed_df = df_concat.select(F.create_map(F.lit("CxName"), "CxName").alias("kv_map"), "Acct") \ .select("Acct", F.explode("kv_map")) \ .groupBy("key") \ .pivot("Acct") \ .agg(F.first("value")) \ .withColumnRenamed("key", "Column") # 查看结果 transposed_df.show()
最终效果说明
转置后的DataFrame中,Column列值为CxName,对应Acct列的值为按位置排序后的逗号分隔字符串(如John Doe, Jane Doe),与目标效果一致。
内容的提问来源于stack exchange,提问作者Dan A
相关产品推荐
相关产品推荐

