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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 03:16:24