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

