PySpark实现两列元素无前置重复的顺序去重处理
PySpark筛选两列均无重复历史记录的行
你需要筛选PySpark DataFrame的行,要求最终结果中Column1和Column2各自都没有重复值,且仅保留原顺序中首次满足“该列值未在之前保留行中出现”的行(任意一列值已在保留行中出现过的行直接忽略)。
输入DataFrame
| Column1 | Column2 |
|---|---|
| A | 1 |
| B | 2 |
| A | 3 |
| C | 3 |
| C | 4 |
| C | 4 |
| D | 4 |
| E | 5 |
| F | 4 |
| G | 7 |
| D | 8 |
| H | 9 |
| I | 9 |
| H | 10 |
| I | 10 |
期望结果DataFrame
| Column1 | Column2 |
|---|---|
| A | 1 |
| B | 2 |
| C | 3 |
| D | 4 |
| E | 5 |
| G | 7 |
| H | 9 |
| I | 10 |
实现思路
由于需要全局跟踪已保留的Column1和Column2值(非分组或滑动窗口场景),PySpark DataFrame的常规窗口函数难以直接实现这种状态累积逻辑。因此可以借助RDD的迭代处理能力,按顺序遍历每行并维护已出现值的集合,仅保留两列值均未出现过的行。
代码实现
- 初始化SparkSession并创建输入DataFrame
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType spark = SparkSession.builder.appName("UniqueColumnRows").getOrCreate() # 输入数据 data = [ ("A", 1), ("B", 2), ("A", 3), ("C", 3), ("C", 4), ("C", 4), ("D", 4), ("E", 5), ("F", 4), ("G", 7), ("D", 8), ("H", 9), ("I", 9), ("H", 10), ("I", 10) ] # 定义Schema schema = StructType([ StructField("Column1", StringType(), True), StructField("Column2", IntegerType(), True) ]) df = spark.createDataFrame(data, schema)
- 添加行号以保留原顺序
from pyspark.sql.functions import monotonically_increasing_id # 添加唯一行号,确保后续处理顺序与输入一致 df_with_id = df.withColumn("row_id", monotonically_increasing_id()).orderBy("row_id")
- 转换为RDD并筛选符合条件的行
def filter_unique_rows(iterator): # 维护已保留的Column1和Column2值集合 seen_col1 = set() seen_col2 = set() for row in iterator: col1_val = row.Column1 col2_val = row.Column2 # 仅当两列值均未出现过时,保留该行并更新集合 if col1_val not in seen_col1 and col2_val not in seen_col2: yield row seen_col1.add(col1_val) seen_col2.add(col2_val) # 转换为RDD处理后转回DataFrame,移除行号列 result_rdd = df_with_id.rdd.mapPartitions(filter_unique_rows) result_df = spark.createDataFrame(result_rdd, df_with_id.schema).drop("row_id")
- 查看结果
result_df.show()
运行后即可得到符合要求的结果。
内容的提问来源于stack exchange,提问作者mouli lee
相关产品推荐
相关产品推荐

