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

Spark DataFrame调用duplicated报错,如何转为Pandas DataFrame?

问题解答

报错原因

Spark DataFrame与Pandas DataFrame属于不同的计算框架API,duplicated()是Pandas DataFrame专属方法,Spark DataFrame并未实现该属性,直接调用自然会抛出AttributeError。

正确查找Spark DataFrame重复行的方法

方法1:获取所有重复行(包含每组重复的全部记录)

from pyspark.sql import functions as F

# 按指定列分组统计行数
count_group = df.groupBy("variable1", "variable2").count()
# 关联原表并筛选出行数>1的重复记录
duplicate_rows = df.join(count_group, on=["variable1", "variable2"], how="inner") \
                   .filter(count_group["count"] > 1) \
                   .drop("count")

duplicate_rows.show()

方法2:实现类似Pandas duplicated(keep='first')的效果

from pyspark.sql.window import Window

# 按指定列分区,指定排序字段
window_spec = Window.partitionBy("variable1", "variable2").orderBy("variable3")
# 给每组记录添加行号,行号>1的即为重复行
duplicate_rows = df.withColumn("row_num", F.row_number().over(window_spec)) \
                   .filter(F.col("row_num") > 1) \
                   .drop("row_num")

duplicate_rows.show()

Spark DataFrame转Pandas DataFrame的方法

直接调用Spark DataFrame的toPandas()方法即可:

pandas_df = df.toPandas()

注意:该操作会将分布式存储的Spark数据拉取到单机内存中,仅适合小数据集转换,大数据量场景下容易触发内存溢出。

内容的提问来源于stack exchange,提问作者Kenzo-San

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 23:15:57