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

如何在PySpark中提取仅含Type='A'的ID对应的行?

解决PySpark提取仅含唯一Type="A"的ID行的问题

要找出那些**所有记录Type都只有"A"**的ID对应的行,我这里有两种实用的PySpark实现方法,都能精准得到你想要的结果:

方法一:分组聚合+过滤+关联

这个方法思路很直观,先统计每个ID的Type分布情况,筛选出符合条件的ID后再关联回原表获取完整数据:

from pyspark.sql.functions import collect_set, size, col

# 第一步:按ID分组,收集每个ID的所有唯一Type,同时统计Type的种类数
id_type_summary = df.groupBy("ID")\
    .agg(collect_set("Type").alias("unique_types"), size(collect_set("Type")).alias("type_count"))

# 第二步:筛选出只有Type="A"的ID(即unique_types是["A"]且种类数为1)
valid_ids = id_type_summary.filter((col("unique_types") == ["A"]) & (col("type_count") == 1))\
    .select("ID")

# 第三步:关联原表,拿到这些ID对应的行数据
final_result = df.join(valid_ids, on="ID", how="inner")
final_result.show()

方法二:窗口函数(更高效)

用窗口函数可以省去额外的关联操作,直接在原数据集上计算每个ID的Type分布,然后过滤出符合条件的行,数据量大时性能更优:

from pyspark.sql.window import Window
from pyspark.sql.functions import collect_set, size, col

# 定义窗口规则:以ID为分组依据
id_window = Window.partitionBy("ID")

# 第一步:为每一行添加所属ID的唯一Type集合和Type种类数字段
df_with_metadata = df.withColumn("unique_types", collect_set("Type").over(id_window))\
    .withColumn("type_count", size(col("unique_types")).over(id_window))

# 第二步:筛选出符合条件的行,同时移除临时字段
final_result = df_with_metadata.filter((col("unique_types") == ["A"]) & (col("type_count") == 1))\
    .drop("unique_types", "type_count")
final_result.show()

运行任意一种方法,都会得到你期望的输出:

+---+----+-----+
| ID|Type|Value|
+---+----+-----+
|ID2|   A|   12|
|ID4|   A|   10|
+---+----+-----+

小补充

  • collect_set用来收集每个ID下不重复的Type值,这样如果一个ID有多个不同Type,集合大小就会大于1;
  • 如果你的数据里同一个ID不会重复出现Type="A"的记录,也可以用countDistinct("Type")替代size(collect_set("Type")),逻辑是完全一致的;
  • 窗口函数的方式避免了join操作,减少了数据shuffle的开销,在大数据场景下更推荐使用。

内容的提问来源于stack exchange,提问作者Hardik Gupta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:56:12