PySpark:如何筛选Rank1与Rank2中Name不同的ID数据
PySpark实现按ID过滤Rank1和Rank2名称不同的数据
需求说明
需要筛选出满足以下条件的ID对应的所有记录:
- 该ID同时存在Rank=1和Rank=2的记录
- Rank=1和Rank=2对应的Name字段值不相同
排除仅存在单个Rank,或两个Rank的Name相同的ID。
实现方案(窗口函数版)
使用窗口函数按ID分区,计算每个ID的Rank总数、Rank1的Name和Rank2的Name,再通过过滤条件筛选目标数据,相比groupBy+join更简洁高效。
1. 构建测试DataFrame
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("FilterRankData").getOrCreate() # 创建测试数据 data = [ (1, "A", 10, 1), (1, "B", 20, 2), (2, "C", 10, 1), (2, "C", 12, 2), (3, "D", 11, 1), (4, "E", 12, 1), (4, "E", 13, 2) ] df = spark.createDataFrame(data, schema=["ID", "Name", "Score", "Rank"]) df.show()
2. 定义窗口并计算辅助字段
# 按ID分区的窗口 window = Window.partitionBy("ID") # 添加辅助字段:Rank总数、Rank1的Name、Rank2的Name df_with_aux = df.withColumn("rank_count", F.count("Rank").over(window)) \ .withColumn("rank1_name", F.when(F.col("Rank") == 1, F.col("Name")).over(window)) \ .withColumn("rank2_name", F.when(F.col("Rank") == 2, F.col("Name")).over(window)) # 提取每个ID下Rank1和Rank2的唯一Name值 df_with_aux = df_with_aux.withColumn("rank1_name", F.first("rank1_name").over(window)) \ .withColumn("rank2_name", F.first("rank2_name").over(window))
3. 过滤目标数据
# 过滤条件:Rank总数为2,且Rank1和Rank2的Name不同 result_df = df_with_aux.filter( (F.col("rank_count") == 2) & (F.col("rank1_name") != F.col("rank2_name")) ).drop("rank_count", "rank1_name", "rank2_name") # 展示结果 result_df.show()
执行结果
+---+----+-----+----+ | ID|Name|Score|Rank| +---+----+-----+----+ | 1| A| 10| 1| | 1| B| 20| 2| +---+----+-----+----+
方案说明
- 窗口函数
partitionBy("ID")确保每个ID的记录被单独处理 count("Rank").over(window)统计每个ID的Rank数量,确保同时存在Rank1和Rank2- 通过
when函数结合窗口的first聚合,获取每个ID下Rank1和Rank2对应的Name值 - 最后通过过滤条件筛选出符合要求的记录,无需额外的join操作,逻辑更直接高效
内容的提问来源于stack exchange,提问作者ahajib
相关产品推荐
相关产品推荐

