如何在PySpark中分组过滤DataFrame:筛选无'In'状态支持的链路
PySpark筛选所有关联Support状态均不为'In'的Link Name
问题描述
给定如下结构的DataFrame,需要找出所有**关联的Support的Status均不为'In'**的Link Name,最终结果仅保留Link3:
| Link Name | Support | Status |
|---|---|---|
| Link1 | Support1 | In |
| Link1 | Support2 | In |
| Link1 | Support3 | Out |
| Link2 | Support4 | In |
| Link2 | Support5 | In |
| Link3 | Support6 | Out |
| Link3 | Support7 | Out |
解决方案
以下提供三种可行的实现方式,可根据数据规模和需求选择:
方法1:分组聚合统计过滤
通过分组统计每个Link Name下状态为'In'的记录数,筛选计数为0的项:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, sum, when # 初始化SparkSession spark = SparkSession.builder.appName("filter_valid_links").getOrCreate() # 创建示例DataFrame data = [ ("Link1", "Support1", "In"), ("Link1", "Support2", "In"), ("Link1", "Support3", "Out"), ("Link2", "Support4", "In"), ("Link2", "Support5", "In"), ("Link3", "Support6", "Out"), ("Link3", "Support7", "Out") ] df = spark.createDataFrame(data, ["Link Name", "Support", "Status"]) # 分组统计并过滤 aggregated_df = df.groupBy("Link Name") \ .agg(sum(when(col("Status") == "In", 1).otherwise(0)).alias("in_record_count")) result = aggregated_df.filter(col("in_record_count") == 0).select("Link Name") result.show()
方法2:窗口函数标记过滤
利用窗口函数标记每个Link Name是否存在'In'状态的记录,再筛选标记为0的项:
from pyspark.sql.window import Window from pyspark.sql.functions import max, when window_spec = Window.partitionBy("Link Name") # 添加标记列:存在'In'则为1,否则为0 df_with_flag = df.withColumn( "has_in_status", max(when(col("Status") == "In", 1).otherwise(0)).over(window_spec) ) # 过滤并去重 result = df_with_flag.filter(col("has_in_status") == 0).select("Link Name").distinct() result.show()
方法3:反选排除法
先找出所有包含'In'状态的Link Name,再从原数据中排除这些项:
# 找出所有存在'In'状态的Link Name invalid_links = df.filter(col("Status") == "In").select("Link Name").distinct() # 左反连接排除无效项,再去重 result = df.join(invalid_links, on="Link Name", how="left_anti") \ .select("Link Name").distinct() result.show()
输出结果
三种方法最终都会输出:
+---------+ |Link Name| +---------+ | Link3| +---------+
内容的提问来源于stack exchange,提问作者SaurabhShelar
相关产品推荐
相关产品推荐

