PySpark实现两列值相同组合(含逆序)的行去重
邮编无序对去重实现方案
需求核心:将ZIP1、ZIP2顺序相反的行视为重复数据,每组重复数据仅保留任意1行,支持百万行级数据处理。
示例输入数据:
import pandas as pd from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() df = pd.DataFrame({'ZIP1': ['50069', '50069', '50704', '50704', '52403', '52403'], 'ZIP2': ['50704', '52403', '50069', '52403', '50069', '50704'], 'STATE': ['IA', 'IA', 'IA', 'IA', 'IA', 'IA'], 'REGION': ['MIDWEST', 'MIDWEST', 'MIDWEST', 'MIDWEST', 'MIDWEST', 'MIDWEST'] } ) sdf = spark.createDataFrame(df)
处理后输出示例(保留行可根据规则调整,无顺序要求):
+-----+-----+-----+-------+ | ZIP1| ZIP2|STATE| REGION| +-----+-----+-----+-------+ |50069|50704| IA|MIDWEST| |50069|52403| IA|MIDWEST| |50704|52403| IA|MIDWEST| +-----+-----+-----+-------+
PySpark 实现(推荐,适配百万级以上大规模数据)
核心思路:使用内置array_sort函数对每行的两个邮编组成的数组做排序,生成和顺序无关的唯一组合键,基于该键做去重,全程使用Spark内置Catalyst优化的函数,性能远高于自定义UDF。
- 基础高性能版本(同组合下其他字段值一致时使用,性能最优)
from pyspark.sql import functions as F result_sdf = sdf.withColumn( "zip_pair", F.array_sort(F.array("ZIP1", "ZIP2")) # 注意:如果同邮编对下STATE/REGION可能不同,可根据业务调整去重字段 ).dropDuplicates(["zip_pair", "STATE", "REGION"]).drop("zip_pair") result_sdf.show()
- 可控保留规则版本(需要指定保留某类行时使用)
如果需要自定义每组重复数据保留的行(比如保留ZIP1更小的行),可以用窗口函数实现:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 按排序后的邮编对、公共维度字段分组,自定义排序规则选要保留的行 win = Window.partitionBy( F.array_sort(F.array("ZIP1", "ZIP2")), "STATE", "REGION" ).orderBy(F.asc("ZIP1")) # 可修改排序逻辑,比如按入库时间倒序保留最新行 result_sdf = sdf.withColumn("rn", F.row_number().over(win))\ .filter(F.col("rn") == 1)\ .drop("rn")
性能调优提示:处理百万行以上数据时,可根据集群资源调整spark.sql.shuffle.partitions参数(一般设置为CPU核数的2~3倍),减少shuffle开销。
Pandas 实现(适合中小规模数据集)
核心逻辑和Spark一致,先生成顺序无关的邮编对键,再基于键去重。
- 易读版本
# 对每行两个邮编排序生成元组作为去重键 df["pair_key"] = df.apply(lambda x: tuple(sorted([x["ZIP1"], x["ZIP2"]])), axis=1) # keep参数可选first/last,决定保留第一行还是最后一行 result_df = df.drop_duplicates( subset=["pair_key", "STATE", "REGION"], keep="first" ).drop(columns=["pair_key"]).reset_index(drop=True)
- 高性能版本(避免apply循环,适配百万行Pandas数据)
import numpy as np # 用numpy向量化排序替代逐行apply,速度提升5~10倍 zip_mat = df[["ZIP1", "ZIP2"]].to_numpy() zip_mat.sort(axis=1) df["pair_key"] = zip_mat[:, 0] + "_" + zip_mat[:, 1] result_df = df.drop_duplicates( subset=["pair_key", "STATE", "REGION"], keep="first" ).drop(columns=["pair_key"]).reset_index(drop=True)
内容的提问来源于stack exchange,提问作者zesla
相关产品推荐
相关产品推荐

