在PySpark中实现Pivot聚合:用对应行answer填充列
PySpark实现Pivot聚合转换列
构建示例输入数据
先创建与示例匹配的DataFrame:
from pyspark.sql import SparkSession from pyspark.sql.functions import first spark = SparkSession.builder.appName("PivotExample").getOrCreate() data = [ (1, "quest_1", "Good"), (1, "quest_2", "Bad"), (2, "quest_1", "Bad"), (2, "quest_2", "Good"), (2, "quest_3", "Quite Good") ] df = spark.createDataFrame(data, ["id", "question", "answer"])
执行Pivot聚合操作
通过groupBy按id分组,用pivot将question的取值转为新列,最后用first聚合获取对应answer值(由于每个id+question组合仅对应一条记录,first可准确拿到唯一的answer):
pivoted_df = df.groupBy("id").pivot("question").agg(first("answer"))
查看输出结果
运行show()方法即可得到期望格式的结果:
pivoted_df.show()
输出内容:
+---+-------+-------+-----------+ | id|quest_1|quest_2| quest_3| +---+-------+-------+-----------+ | 1| Good| Bad| null| | 2| Bad| Good|Quite Good| +---+-------+-------+-----------+
关键说明
pivot("question"):自动识别question列的所有唯一取值,并将这些值作为新的列名。- 聚合函数选择:如果数据中
id+question存在多条重复记录,可根据需求替换为max、min或自定义聚合函数;若为唯一记录,first/last均适用。
内容的提问来源于stack exchange,提问作者Jresearcher
相关产品推荐
相关产品推荐

