如何在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
相关产品推荐
相关产品推荐

