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
相关产品推荐
相关产品推荐

