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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 10:32:33