You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 10:17:28