Spark DataFrame按相同键合并多行非空值为单行的实现方法
大规模Spark DataFrame按组合并非空值单行方案
核心思路
超大数据量场景下优先使用Spark原生内置聚合函数,避免UDF、自定义聚合、数组构造这类高开销逻辑,依托Catalyst优化器的全阶段代码生成和Map端预聚合能力,性能比常规自定义方案高一个数量级,完全适配TB级数据规模。
具体实现
逻辑是按指定分组键聚合,对每个非分组字段直接取分组内第一个非空值,用到的first函数支持ignorenulls=True参数,会自动跳过空值取第一个有效值,全程无额外内存开销。
from pyspark.sql import functions as F # 替换为你实际用来分组的键列名 group_columns = ["你的分组键1", "你的分组键2"] # 自动构造所有非分组列的聚合规则 agg_rules = [ F.first(col_name, ignorenulls=True).alias(col_name) for col_name in df.columns if col_name not in group_columns ] # 执行聚合得到结果 result = df.groupBy(*group_columns).agg(*agg_rules)
超大规模数据优化点
- 如果同组数据存在多个不同非空值,需要取最新值/最大值的话,把
first替换为last/max即可,这两个函数同样支持ignorenulls=True参数,性能无损耗 - 遇到数据倾斜场景(单个分组键对应数据量远超其他分组),可以加两阶段聚合优化:第一阶段给分组键拼接1~20范围内的随机前缀,做局部预聚合,第二阶段去掉随机前缀做全局聚合,避免单节点OOM
- 禁止使用
collect_list+UDF遍历数组合并非空值的方案:这类方案会把同组所有数据拉到Executor内存构造数组,不仅序列化开销大,数据量稍大就会直接内存溢出,shuffle传输量也远高于原生聚合方案
内容的提问来源于stack exchange,提问作者rahul kishore
相关产品推荐
相关产品推荐

