PySpark:如何基于窗口函数分组结果进行二次数据筛选?
PySpark 分组筛选解决方案
实现思路
核心是先通过窗口函数获取每个firstID+secondID分组内的所有type集合,再根据给定规则匹配需要保留的目标type,最后筛选出对应记录。
完整代码示例
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义分组窗口(仅按firstID、secondID分区,无需排序) group_window = Window.partitionBy("firstID", "secondID") # 1. 计算每个分组的type集合,确定目标保留的type df_with_target = df.withColumn( "type_set", F.collect_set("type").over(group_window) ).withColumn( "target_type", F.when(F.array_contains(F.col("type_set"), "type1") & F.array_contains(F.col("type_set"), "type2") & F.array_contains(F.col("type_set"), "type3"), "type3") .when(F.array_contains(F.col("type_set"), "type1") & F.array_contains(F.col("type_set"), "type2"), "type2") .when(F.array_contains(F.col("type_set"), "type1") & F.array_contains(F.col("type_set"), "type3"), "type3") .otherwise(F.col("type")) ) # 2. 筛选出type与目标type匹配的记录 result_df = df_with_target.filter(F.col("type") == F.col("target_type")).drop("type_set", "target_type") # 查看结果 result_df.show(truncate=False)
代码说明
collect_set("type").over(group_window):收集每个分组内的所有不重复type,生成集合列type_set,用于判断分组内的type组合。target_type列逻辑:严格按照给定的4条规则依次判断,匹配到对应条件后返回要保留的type;单条记录时直接返回自身type。- 筛选逻辑:只保留
type等于target_type的行,每个分组最终只会留下符合规则的一条记录。
补充场景处理
如果分组内存在多条同类型的记录(比如同一分组有两条type3),可以结合你之前定义的窗口,取最新的一条记录:
# 复用你定义的窗口 windowSpec = Window.partitionBy("firstID","secondID").orderBy(F.col("timestamp").desc()) # 在筛选前添加行号,取每个目标type的第一条(最新记录) result_df = df_with_target.filter(F.col("type") == F.col("target_type")) \ .withColumn("row_num", F.row_number().over(windowSpec)) \ .filter(F.col("row_num") == 1) \ .drop("type_set", "target_type", "row_num")
内容的提问来源于stack exchange,提问作者peace
相关产品推荐
相关产品推荐

