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
相关产品推荐
相关产品推荐

