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

Azure Databricks中PySpark去除主键重复行的实现方法

PySpark 按主键去重(保留首行或随机行)

针对你在Azure Databricks中的需求,这里提供两种实用的PySpark脚本实现,分别对应保留首行和随机保留一行的场景:

1. 先构建示例DataFrame

先创建你提供的测试数据,方便验证效果:

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, rand, monotonically_increasing_id

# Databricks环境中无需手动初始化SparkSession,可直接跳过这行
spark = SparkSession.builder.appName("ID_Deduplication").getOrCreate()

# 构造测试数据
data = [
    (1, "Mike", 18, "USA"),
    (2, "Lily", 25, "Japan"),
    (2, "Lily", 22, "Japan"),
    (3, "John", 23, "China"),
    (4, "Snow", 36, "Korea"),
    (5, "Tom", 28, "Austria"),
    (5, "Cindy", 28, "Italy")
]

df = spark.createDataFrame(data, ["id", "Name", "Age", "Country"])
df.show()

2. 保留每个id的首行(严格匹配输入顺序)

由于Spark是分布式计算框架,默认行顺序不固定,要确保保留输入时的第一行,需要先添加自增序号标记顺序,再通过窗口函数筛选:

# 添加自增序号列,标记数据输入顺序
df_with_seq = df.withColumn("seq", monotonically_increasing_id())

# 按id分组,取序号最小的行(即输入首行)
window_spec = Window.partitionBy("id").orderBy("seq")
df_first_row = df_with_seq.withColumn("row_num", row_number().over(window_spec)) \
    .filter("row_num = 1") \
    .drop("seq", "row_num")

df_first_row.show()

执行后会得到你要的结果一:保留id=2的Age25行、id=5的Tom行。

3. 随机保留每个id的一行

如果不需要固定保留首行,仅需随机去重,有两种简洁方式:

方式A:直接使用dropDuplicates

这是最简便的方法,Spark会自动为每个id保留任意一行:

df_rand_simple = df.dropDuplicates(["id"])
df_rand_simple.show()

方式B:用随机排序的窗口函数

通过rand()函数随机排序后取第一行,可控性更强:

window_spec_rand = Window.partitionBy("id").orderBy(rand())
df_rand = df.withColumn("row_num", row_number().over(window_spec_rand)) \
    .filter("row_num = 1") \
    .drop("row_num")

df_rand.show()

这两种方式都可能得到你要的结果二(保留id=2的Age22行、id=5的Cindy行),每次执行结果可能略有不同。

内容的提问来源于stack exchange,提问作者Emma

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 01:23:29