PySpark中如何结合条件过滤与GroupBy实现聚合去重?
问题描述
现有Spark DataFrame df_A 数据如下:
+-------------+-------------+--------+-----------+--------+ | id| Type|id_count| Value1| Value2| +-------------+-------------+--------+-----------+--------+ | 18| AAA| 2| null| null| | 18| BBB| 2| null| null| | 16| BBB| 2| null| null| | 16| CCC| 2| null| null| | 17| CCC| 1| null| null| +-------------+-------------+--------+-----------+--------+
需求:当id_count等于2且该id下无Type="AAA"的记录时,对该id的Value1求和,同时合并去重后的Type字段;其余情况保留原始行数据。最终期望结果如下:
+-------------+-------------+--------+-----------+--------+ | id| Type|id_count| Value1| Value2| +-------------+-------------+--------+-----------+--------+ | 18| AAA| 2| null| null| | 18| BBB| 2| null| null| | 16| CCC+BBB| 2| null| null| | 17| CCC| 1| null| null| +-------------+-------------+--------+-----------+--------+
常规GroupBy求和操作(df_B = df_A.groupby(id).sum('Value1'))无法满足条件筛选与字段合并的组合需求,需结合窗口函数、条件判断实现。
解决方案
通过分组标记+分情况处理+结果合并的方式实现,具体代码如下(基于PySpark):
from pyspark.sql import functions as F from pyspark.sql.window import Window # 步骤1:标记每个id是否需要合并(id_count=2且该id下无AAA类型) window_spec = Window.partitionBy("id") df_marked = df_A.withColumn( "need_merge", F.when( (F.col("id_count") == 2) & (F.count(F.when(F.col("Type") == "AAA", 1)).over(window_spec) == 0), True ).otherwise(False) ) # 步骤2:处理需要合并的分组:合并Type、求和Value1 merged_df = df_marked.filter(F.col("need_merge") == True) \ .groupBy("id", "id_count") \ .agg( F.concat_ws("+", F.collect_set("Type")).alias("Type"), F.sum(F.col("Value1")).alias("Value1"), F.first("Value2").alias("Value2") # Value2取第一个值,可根据需求调整 ) # 步骤3:处理不需要合并的分组:保留原始行数据 unmerged_df = df_marked.filter(F.col("need_merge") == False) \ .select("id", "Type", "id_count", "Value1", "Value2") # 步骤4:合并两个结果集并按id、Type排序 final_df = merged_df.unionByName(unmerged_df) \ .orderBy("id", "Type") final_df.show()
代码说明:
- 分组标记:用窗口函数统计每个
id下Type=AAA的记录数,结合id_count=2的条件,标记出需要合并的id组。 - 合并目标分组:对标记为需要合并的组,用
collect_set去重后通过concat_ws合并Type字段,用sum计算Value1的总和;Value2根据实际需求选择保留方式(示例中取第一个值)。 - 保留原数据:对不需要合并的组直接保留原始行。
- 结果合并:将合并后的分组数据与原始未合并数据合并,排序后得到最终结果。
若需求中“对Value1求和并去重”指的是对Value1的非重复值求和,可将F.sum(F.col("Value1"))替换为F.sum(F.expr("array_distinct(collect_list(Value1))"))(需注意空值处理)。
内容的提问来源于stack exchange,提问作者otk
相关产品推荐
相关产品推荐

