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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 05:50:31