如何移除Block列中对应多个不同ID值的所有相关行?
如何移除对应多个不同ID的Block的所有行
需求说明
需要移除满足以下条件的所有行:当某个Block值对应至少两个不同的ID时,删除所有包含该Block值的行。
示例数据
原始数据:
| ID | Block |
|---|---|
| 1 | A |
| 1 | C |
| 1 | C |
| 3 | A |
| 3 | B |
期望输出:
| ID | Block |
|---|---|
| 1 | C |
| 1 | C |
| 3 | B |
解决方案(PySpark)
这里提供两种实现方式,均能达到需求:
方法1:分组统计+关联过滤
- 先统计每个Block对应的唯一ID数量
- 筛选出仅对应单个ID的Block
- 将原始数据与这些合法Block关联,得到结果
from pyspark.sql import SparkSession from pyspark.sql.functions import countDistinct # 初始化Spark会话 spark = SparkSession.builder.appName("filter_block_by_id").getOrCreate() # 构造原始数据 data = [(1, "A"), (1, "C"), (1, "C"), (3, "A"), (3, "B")] df = spark.createDataFrame(data, ["ID", "Block"]) # 统计每个Block的唯一ID数量 block_id_stats = df.groupBy("Block").agg(countDistinct("ID").alias("distinct_id_count")) # 筛选出仅对应1个ID的Block valid_blocks = block_id_stats.filter(block_id_stats.distinct_id_count == 1).select("Block") # 过滤原始数据,保留合法Block的行 result_df = df.join(valid_blocks, on="Block", how="inner") result_df.show()
方法2:窗口函数直接过滤
使用窗口函数给每行标记其所属Block的唯一ID数量,再直接筛选符合条件的行,无需额外关联:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import countDistinct spark = SparkSession.builder.appName("filter_block_by_id").getOrCreate() data = [(1, "A"), (1, "C"), (1, "C"), (3, "A"), (3, "B")] df = spark.createDataFrame(data, ["ID", "Block"]) # 定义窗口:按Block分组 window_spec = Window.partitionBy("Block") # 给每行添加该Block对应的唯一ID数量 df_with_stats = df.withColumn("distinct_id_count", countDistinct("ID").over(window_spec)) # 筛选出唯一ID数量为1的行,移除辅助列 result_df = df_with_stats.filter(df_with_stats.distinct_id_count == 1).drop("distinct_id_count") result_df.show()
关键逻辑说明
两种方法的核心都是:识别出那些仅对应单个ID的Block,然后保留这些Block的所有行。对于Block值对应多个不同ID的情况(比如示例中的Block A),所有相关行都会被过滤掉。
内容的提问来源于stack exchange,提问作者john
相关产品推荐
相关产品推荐

