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

PySpark按id、cod分组时如何去重并过滤含非空值分组的空flag行

PySpark 分组过滤DataFrame实现

需求说明

现有PySpark DataFrame,需按id、cod列分组,基于flag列执行过滤,规则如下:

  • 分组内不存在flag值不为None的行:仅保留1条唯一行
  • 分组内存在flag值非None的行:先移除分组内flag为None的行,再对剩余重复行去重

初始测试代码

import pyspark
from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import col, row_number,max

spark = SparkSession.builder.appName('Vazio').getOrCreate()

data = [('1', 10, 'A'),
('1', 10, 'A'),
('1', 10, None),
('1', 15, 'A'),
('1', 15, None),
('2', 11, 'A'),
('2', 11, 'C'),
('2', 12, 'B'),
('2', 12, 'B'),
('2', 12, 'C'),
('2', 12, 'C'),
('2', 13, None),
('3', 14, None),
('3', 14, None),
('3', 15, None),
('4', 21, 'A'),
('4', 21, 'B'),
('4', 21, 'C'),
('4', 21, 'C')]

df = spark.createDataFrame(data=data, schema = ['id', 'cod','flag'])
df.show()

目标输出

预期得到的DataFrame如下:

+---+---+----+
| id|cod|flag|
+---+---+----+
|  1| 10|   A|
|  1| 15|   A|
|  2| 11|   A|
|  2| 11|   C|
|  2| 12|   B|
|  2| 12|   C|
|  2| 13|null|
|  3| 14|null|
|  3| 15|null|
|  4| 21|   A|
|  4| 21|   C|
+---+---+----+

实现方案

通过窗口函数标记每个分组的非空flag存在状态,再按规则过滤后去重即可,代码如下:

from pyspark.sql.functions import count, when, lit

# 定义id、cod维度的分组窗口
group_window = Window.partitionBy("id", "cod")

# 为每行添加分组标记:判断分组内是否存在非空flag
df_process = df.withColumn(
    "has_non_null_flag",
    count(when(col("flag").isNotNull(), lit(1))).over(group_window) > 0
)

# 按规则过滤行
result = df_process.filter(
    # 分组存在非空flag时剔除null行,无有效flag时保留null行
    when(col("has_non_null_flag"), col("flag").isNotNull())
    .otherwise(col("flag").isNull())
).select("id", "cod", "flag").dropDuplicates()

result.show()

针对输出中id=4,cod=21分组缺失flag=B行的问题,是原始样例输出的笔误,按照给定规则该行属于非空、无重复的有效行,运行上述代码会正常返回,若有特殊业务规则剔除B值可在过滤条件中补充对应逻辑即可。


内容的提问来源于stack exchange,提问作者Pedro Henrique

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 23:21:24