如何在PySpark、Pandas中提取指定列重复行并从原DataFrame删除?
Pandas 实现方法
- 步骤1:标记所有基于
Articlenbr和amount列的重复行(keep=False会标记所有重复出现的行,而非仅排除第一/最后一行外的重复行) - 步骤2:提取重复行到新DataFrame
- 步骤3:从原DataFrame中移除重复行,保留唯一行
import pandas as pd # 示例DataFrame data = { 'Articlenbr': ['A001', 'A001', 'B002', 'C003'], 'amount': [100, 100, 200, 300], 'other_col': ['x', 'y', 'z', 'w'] } df = pd.DataFrame(data) # 提取重复行 df_duplicates = df[df.duplicated(subset=['Articlenbr', 'amount'], keep=False)] # 原DataFrame保留非重复行 df = df[~df.duplicated(subset=['Articlenbr', 'amount'], keep=False)]
PySpark 实现方法
方法1:分组计数关联
from pyspark.sql import SparkSession from pyspark.sql.functions import count # 初始化SparkSession spark = SparkSession.builder.appName("DuplicateHandling").getOrCreate() # 示例DataFrame data = [ ("A001", 100, "x"), ("A001", 100, "y"), ("B002", 200, "z"), ("C003", 300, "w") ] df = spark.createDataFrame(data, ["Articlenbr", "amount", "other_col"]) # 分组统计每个(Articlenbr, amount)的出现次数 count_df = df.groupBy("Articlenbr", "amount").agg(count("*").alias("cnt")) # 关联原DataFrame,标记重复行 df_with_cnt = df.join(count_df, on=["Articlenbr", "amount"], how="left") # 提取重复行(出现次数>1) df_duplicates = df_with_cnt.filter(df_with_cnt.cnt > 1).drop("cnt") # 原DataFrame保留非重复行(出现次数=1) df = df_with_cnt.filter(df_with_cnt.cnt == 1).drop("cnt")
方法2:窗口函数统计
from pyspark.sql.window import Window from pyspark.sql.functions import count window_spec = Window.partitionBy("Articlenbr", "amount") df_with_cnt = df.withColumn("cnt", count("*").over(window_spec)) # 提取重复行 df_duplicates = df_with_cnt.filter("cnt > 1").drop("cnt") # 原DataFrame保留非重复行 df = df_with_cnt.filter("cnt == 1").drop("cnt")
内容的提问来源于stack exchange,提问作者user19814628
相关产品推荐
相关产品推荐

