如何在PySpark中根据另一列将列转换为列表
PySpark按分组聚合列值为列表的解决方案
你需要的是按Column A分组,将对应Column B的所有值聚合为列表,map()方法不适合这个场景——因为map是逐行处理,无法实现跨行的分组合并。正确的做法是使用Spark内置的分组聚合函数collect_list,具体步骤如下:
1. 导入依赖并创建示例DataFrame
from pyspark.sql import SparkSession from pyspark.sql.functions import collect_list # 初始化Spark会话 spark = SparkSession.builder.appName("GroupCollectList").getOrCreate() # 构造你的示例数据 data = [(123, "abc"), (123, "def"), (456, "klm"), (789, "nop"), (789, "qrst")] df = spark.createDataFrame(data, ["Column A", "Column B"])
2. 执行分组聚合
通过groupBy指定分组字段,再用collect_list聚合对应列的值为列表:
result_df = df.groupBy("Column A").agg(collect_list("Column B").alias("Column B")) # 查看结果 result_df.show(truncate=False)
执行后会得到你期望的输出:
+--------+------------+ |Column A|Column B | +--------+------------+ |123 |[abc, def] | |456 |[klm] | |789 |[nop, qrst] | +--------+------------+
可选:保证列表内元素的顺序
如果需要严格按Column B的排序输出列表,可以先对数据排序再聚合:
# 先排序再分组聚合 sorted_result_df = df.orderBy("Column A", "Column B") \ .groupBy("Column A") \ .agg(collect_list("Column B").alias("Column B")) sorted_result_df.show(truncate=False)
内容的提问来源于stack exchange,提问作者Abhijith Nagarjuna
相关产品推荐
相关产品推荐

