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

PySpark如何转置分组DataFrame?无聚合列时的实现方案

Spark DataFrame 转置实现方案

你可以通过添加分组内序号+ pivot聚合的方式实现需求,具体步骤如下:

核心思路

由于pivot()需要聚合键,我们先给每个group分组内的words添加行号,以此作为聚合依据,再通过pivot将group的取值转为列,最后按需调整第一列内容。

PySpark 实现代码

  1. 初始化示例DataFrame
from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, first, when, lit, col

spark = SparkSession.builder.appName("transpose_demo").getOrCreate()

# 构造原始数据
data = [
    ("one", "word1"), ("one", "word2"), ("one", "word3"),
    ("two", "word4"), ("two", "word5"), ("two", "word6"),
    ("three", "word7"), ("three", "word8"), ("three", "word9")
]
df = spark.createDataFrame(data, ["group", "words"])
  1. 为分组添加行号
# 按group分组,words排序生成行号(排序依据可根据实际需求调整)
window = Window.partitionBy("group").orderBy("words")
df_with_row = df.withColumn("row_num", row_number().over(window))
  1. Pivot转置并调整格式
# 以行号为聚合键,pivot group列,取对应位置的words
pivoted_df = df_with_row.groupBy("row_num").pivot("group").agg(first("words"))

# 按需生成第一列(若不需要可直接跳过此步)
final_df = pivoted_df.withColumn("group", when(col("row_num") == 2, lit("words")).otherwise(lit(""))) \
                     .drop("row_num") \
                     .select("group", "one", "two", "three")

final_df.show()

执行后输出结果:

+------+-----+-----+------+
| group|  one|  two| three|
+------+-----+-----+------+
|      |word1|word4|word7 |
| words|word2|word5|word8 |
|      |word3|word6|word9 |
+------+-----+-----+------+

Scala 实现代码

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

val spark = SparkSession.builder.appName("transpose_demo").getOrCreate()

// 构造原始数据
val data = Seq(
  ("one", "word1"), ("one", "word2"), ("one", "word3"),
  ("two", "word4"), ("two", "word5"), ("two", "word6"),
  ("three", "word7"), ("three", "word8"), ("three", "word9")
)
val df = spark.createDataFrame(data).toDF("group", "words")

// 添加分组行号
val window = Window.partitionBy("group").orderBy("words")
val dfWithRow = df.withColumn("row_num", row_number().over(window))

// Pivot转置并调整格式
val pivotedDf = dfWithRow.groupBy("row_num").pivot("group").agg(first("words"))
val finalDf = pivotedDf.withColumn("group", when(col("row_num") === 2, lit("words")).otherwise(lit("")))
  .drop("row_num")
  .select("group", "one", "two", "three")

finalDf.show()

说明

  • 行号的排序依据可根据实际业务调整(比如不需要按words排序的话,可以用monotonically_increasing_id()或其他方式生成序号)
  • 若不需要第一列,直接去掉withColumn和select中的group相关逻辑即可

内容的提问来源于stack exchange,提问作者Katy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 22:25:31