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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 01:48:16