如何在PySpark中按值实现透视(Pivot)转换?
PySpark 实现DataFrame透视转换的正确方法
问题场景
输入DataFrame:
+----+-----+---+------+----+------+-------+--------+ |year|month|day|new_ts|hour|minute|ts_rank| label| +----+-----+---+------+----+------+-------+--------+ |2022| 1| 1| 13| 13| 24| 1| 7| |2022| 1| 1| 14| 13| 24| 1| 8| |2022| 1| 2| 15| 13| 24| 1| 7| |2022| 1| 2| 16| 13| 44| 7| 8| +----+-----+---+------+----+------+-------+--------+
期望输出:
+----+-----+---+-------+--------+ |year|month|day| 7 | 8| +----+-----+---+-------+--------+ |2022| 1| 1| 13| 14| |2022| 1| 2| 15| 16| +----+-----+---+-------+--------+
Pandas中可通过以下代码实现需求:
df_pivot = df.pivot(index=["year","month","day"], columns="label", values="new_ts").reset_index()
但你在PySpark中使用df.groupBy(["year","month","day"]).pivot("label").value("new_ts")会报错,修正方法如下:
解决方案
PySpark的pivot方法后必须指定聚合函数——因为Spark是分布式计算框架,需要明确分组后的数据聚合规则。你的场景中,每个year-month-day-label组合仅对应一条new_ts,用first()或max()均可得到正确结果:
from pyspark.sql.functions import first df_pivot = df.groupBy(["year", "month", "day"]) \ .pivot("label") \ .agg(first("new_ts")) \ .orderBy("year", "month", "day")
或者用max()实现:
df_pivot = df.groupBy(["year", "month", "day"]) \ .pivot("label") \ .agg(max("new_ts")) \ .orderBy("year", "month", "day")
关键说明
groupBy(["year", "month", "day"]):按年、月、日分组,这部分和你的代码一致。pivot("label"):将label列的唯一值(7、8)转换为新列名。.agg(first("new_ts")):对分组后的new_ts取第一个值,由于你的数据中每个分组+标签组合仅有一条记录,first()能精准匹配对应值。orderBy为可选操作,用于保证输出顺序和预期一致。
报错原因
PySpark的pivot方法返回的是GroupedData对象,无法直接调用value()方法,必须通过agg()指定聚合函数才能生成最终DataFrame——这是PySpark与Pandaspivot方法的核心差异。
内容的提问来源于stack exchange,提问作者Nabih Bawazir
相关产品推荐
相关产品推荐

