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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 13:20:19