如何将PySpark DataFrame按分组列转换为字典列表?
解决PySpark DataFrame分组生成字典列表的问题
你的代码存在两个问题:
- 列名引用错误:原DataFrame的列是
pdct_code,但你写的是pdct_cd - struct参数写法不规范:第二个参数未正确调用列对象或指定别名
正确代码(按需选择)
场景1:将pdct_code重命名为pdct_cd(匹配你预期的字典键名)
from pyspark.sql import functions as F # 假设df是输入DataFrame Df1 = df.groupBy("ID").agg( F.collect_list( F.struct( F.col("accounts"), F.col("pdct_code").alias("pdct_cd") ) ).alias("array_dict") )
场景2:保留原列名pdct_code
from pyspark.sql import functions as F Df1 = df.groupBy("ID").agg( F.collect_list(F.struct("accounts", "pdct_code")).alias("array_dict") )
验证结果
执行Df1.show(truncate=False)会得到结构化输出:
+---+------------------------------------------------+ |ID |array_dict | +---+------------------------------------------------+ |1 |[{100, IN}, {200, CC}] | |2 |[{300, DD}, {400, ZZ}] | |3 |[{500, AA}] | +---+------------------------------------------------+
如果需要转换成纯Python字典列表格式,可通过以下代码实现:
for row in Df1.collect(): dict_list = [item.asDict() for item in row["array_dict"]] print(row["ID"], dict_list)
输出结果:
1 [{'accounts': 100, 'pdct_cd': 'IN'}, {'accounts': 200, 'pdct_cd': 'CC'}] 2 [{'accounts': 300, 'pdct_cd': 'DD'}, {'accounts': 400, 'pdct_cd': 'ZZ'}] 3 [{'accounts': 500, 'pdct_cd': 'AA'}]
内容的提问来源于stack exchange,提问作者Sam Harris
相关产品推荐
相关产品推荐

