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

如何将PySpark DataFrame按指定列分组并转换为字典列表格式

如何将PySpark DataFrame按指定列分组并转换为字典列表格式

嘿,你的思路其实已经踩对重点啦!我来帮你把最后那层窗户纸捅破~

首先得说:你写的代码其实已经非常接近你想要的结果了,只是PySpark的默认显示方式容易给你造成误解。咱们一步步来看:

先明确问题本质

你用F.collect_list(F.struct(*value_cols))得到的并不是“无列名的数组列表”,而是PySpark的Row对象列表——这些Row对象是带字段名的,和你想要的字典结构本质上是等价的,只是显示形式不同而已。

完整可运行的示例代码

咱们先从构建示例数据开始,再一步步实现你要的效果:

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

# 初始化Spark会话
spark = SparkSession.builder.appName("GroupToDictList").getOrCreate()

# 构建你提供的示例DataFrame
data = [
    (1, "a", 0, "metric4", "value6"),
    (1, "a", 0, "metric3", "value7"),
    (1, "b", 1, "metric2", "value5"),
    (2, "b", 0, "metric4", "value8"),
    (2, "c", 0, "metric3", "value9")
]
df = spark.createDataFrame(data, ["A", "B", "C", "D", "E"])

# 执行你写的分组逻辑
value_cols = [col for col in df.columns if col != "A"]
df_test = df.groupBy("A").agg(F.collect_list(F.struct(*value_cols)).alias("ListOfDict"))

查看结果与转换字典

当你运行df_test.show(truncate=False)时,PySpark会把struct显示成Row对象,看起来是这样:

+---+----------------------------------------------------------------------------------------------------+
|A  |ListOfDict                                                                                          |
+---+----------------------------------------------------------------------------------------------------+
|1  |[{a, 0, metric4, value6}, {a, 0, metric3, value7}, {b, 1, metric2, value5}]                         |
|2  |[{b, 0, metric4, value8}, {c, 0, metric3, value9}]                                                  |
+---+----------------------------------------------------------------------------------------------------+

但别担心,这些Row内部是保留了列名的!如果要把它们转换成Python原生的字典列表,只需要在收集结果到Driver端时做个小转换:

# 把DataFrame结果转换成Python字典格式
result = df_test.rdd.map(lambda row: {
    "A": row.A,
    "ListOfDict": [r.asDict() for r in row.ListOfDict]
}).collect()

# 打印看看效果
for item in result:
    print(item)

这时候输出就是你完全想要的格式了:

{'A': 1, 'ListOfDict': [{'B': 'a', 'C': 0, 'D': 'metric4', 'E': 'value6'}, {'B': 'a', 'C': 0, 'D': 'metric3', 'E': 'value7'}, {'B': 'b', 'C': 1, 'D': 'metric2', 'E': 'value5'}]}
{'A': 2, 'ListOfDict': [{'B': 'b', 'C': 0, 'D': 'metric4', 'E': 'value8'}, {'B': 'c', 'C': 0, 'D': 'metric3', 'E': 'value9'}]}

额外选项:转换成JSON字符串列表

如果你希望在DataFrame中直接存储类似字典的JSON字符串,也可以用to_json函数包裹struct:

df_test_json = df.groupBy("A").agg(
    F.collect_list(F.to_json(F.struct(*value_cols))).alias("ListOfDict")
)
df_test_json.show(truncate=False)

输出会是这样的JSON字符串列表:

+---+----------------------------------------------------------------------------------------------------+
|A  |ListOfDict                                                                                          |
+---+----------------------------------------------------------------------------------------------------+
|1  |[{"B":"a","C":0,"D":"metric4","E":"value6"}, {"B":"a","C":0,"D":"metric3","E":"value7"}, {"B":"b","C":1,"D":"metric2","E":"value5"}]|
|2  |[{"B":"b","C":0,"D":"metric4","E":"value8"}, {"B":"c","C":0,"D":"metric3","E":"value9"}]             |
+---+----------------------------------------------------------------------------------------------------+

总结一下

  • 你最初的代码逻辑是对的,collect_list(struct(...))已经保留了列名信息,只是PySpark的默认显示简化了格式;
  • 要得到Python原生字典列表,只需要对Row对象调用.asDict()方法即可;
  • 如果需要JSON格式的字符串,用to_json包裹struct就能实现。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.08 09:54:51