PySpark中按顺序筛选数组列无前置重复元素的行
解决PySpark DataFrame数组元素全局去重筛选问题
需求说明
给定包含两列的PySpark DataFrame:Column1(整数列)、Column2(ArrayType列),需要筛选出**Column2中所有元素均未在之前任意行的Column2中出现过**的行,只要某行Column2存在元素与前置行重复,整行忽略。
解决方案
核心思路是按行顺序维护全局已出现元素集合,逐行检查当前行数组是否与该集合无交集,无交集则保留该行并更新集合。具体实现步骤如下:
1. 导入依赖并构造示例数据
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("FilterUniqueArrayRows").getOrCreate() # 示例数据 data = [ (1, [10, 20]), (2, [30, 40]), (3, [39, 40]), # 因40已在第2行出现,需被过滤 (4, [50, 60]) ] df = spark.createDataFrame(data, ["Column1", "Column2"])
2. 添加行号确保顺序
因为需要按“前置行”逻辑判断,必须先确定行的顺序(这里默认按Column1排序,可根据实际需求调整):
df = df.withColumn("row_num", F.row_number().over(Window.orderBy("Column1")))
3. 收集前置行的所有数组元素
通过窗口函数收集当前行之前所有行的Column2元素,合并为去重后的集合:
# 定义窗口:从第一行到当前行的前一行 window_spec = Window.orderBy("row_num").rowsBetween(Window.unboundedPreceding, Window.currentRow - 1) # 合并前置行的数组元素并去重 df = df.withColumn( "prev_elements", F.flatten(F.collect_list("Column2").over(window_spec)) ).withColumn( "prev_elements_set", F.array_distinct("prev_elements") )
4. 判断当前行是否存在重复元素
通过array_intersect判断当前行Column2与前置元素集合是否有交集,无交集则保留:
df = df.withColumn( "has_duplicate", F.size(F.array_intersect("Column2", "prev_elements_set")) > 0 ) # 筛选无重复的行,清理辅助列 result_df = df.filter(~F.col("has_duplicate")).drop("row_num", "prev_elements", "prev_elements_set", "has_duplicate")
5. 查看结果
result_df.show()
输出结果:
+-------+---------+ |Column1| Column2| +-------+---------+ | 1| [10, 20]| | 2| [30, 40]| | 4| [50, 60]| +-------+---------+
常见问题说明
之前尝试explode数组结合窗口函数未成功,大概率是因为explode后仅处理了单个元素的去重,未还原到整行维度判断——需要确保整行所有元素都未在前置行出现,而非单个元素去重后保留行。
内容的提问来源于stack exchange,提问作者mouli lee
相关产品推荐
相关产品推荐

