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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 20:50:26