Pandas on Spark按A、B分组聚合生成C-D映射列表报错排查
Pandas on Spark分组聚合生成键值映射列表的问题解决
问题描述
使用Pandas on Spark,需按A、B列分组,聚合后返回以C为键、D为值的映射列表。
示例输入
A B C D 0 7 201806851 0006378110 2223982011 1 7 6378110 0006378110 2223982011 2 7 201806851 201806851 20972475011 3 7 6378110 201806851 20972475011
示例输出
A B C 0 7 6378110 [[0006378110, 2223982011], [201806851, 20972475011]] 1 7 201806851 [[0006378110, 2223982011], [201806851, 20972475011]]
出错代码及错误信息
代码第一行触发断言错误:assert len(key) == len(that_column_labels) AssertionError
seed_data["C"] = seed_data[["C", "D"]].to_dict('records') seed_data = (seed_data .groupby(["A", "B"])["C"] .apply(list).reset_index(name="C"))
尝试将C、D列提取到单独DataFrame,转为字典后作为聚合列,仍出现相同错误。
错误原因
Pandas on Spark基于分布式的Spark引擎,而to_dict('records')是生成本地Python字典序列的操作,直接将其赋值给分布式DataFrame的列,会导致本地数据模型和分布式数据模型不兼容,触发内部断言检查失败。
解决方案
使用Spark原生函数构造键值对列表,避免本地与分布式数据的冲突:
from pyspark.sql import functions as F # 构造C、D结构体,按A、B分组聚合为列表 result = seed_data.groupBy("A", "B") \ .agg(F.collect_list(F.struct("C", "D")).alias("C")) \ .toPandas()
若需保留Pandas on Spark DataFrame格式,移除.toPandas()即可:
result = seed_data.groupBy("A", "B") \ .agg(F.collect_list(F.struct("C", "D")).alias("C"))
说明
F.struct("C", "D")将每行的C、D组合为结构体对象F.collect_list()将每组内的结构体收集为列表,匹配示例输出格式- 该方案完全基于Spark分布式操作,不会触发数据模型冲突问题
内容的提问来源于stack exchange,提问作者WorkInProgress
相关产品推荐
相关产品推荐

