PySpark如何转置分组DataFrame?无聚合列时的实现方案
Spark DataFrame 转置实现方案
你可以通过添加分组内序号+ pivot聚合的方式实现需求,具体步骤如下:
核心思路
由于pivot()需要聚合键,我们先给每个group分组内的words添加行号,以此作为聚合依据,再通过pivot将group的取值转为列,最后按需调整第一列内容。
PySpark 实现代码
- 初始化示例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"])
- 为分组添加行号
# 按group分组,words排序生成行号(排序依据可根据实际需求调整) window = Window.partitionBy("group").orderBy("words") df_with_row = df.withColumn("row_num", row_number().over(window))
- 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
相关产品推荐
相关产品推荐

