如何用Scala Spark按主键分组合并DataFrame非空值去重
解决Spark DataFrame字段合并问题
当然可以实现你的需求,核心思路是针对每个字段单独聚合,提取分组内的非空值,而不是收集整个行结构。
具体实现步骤
假设你的DataFrame结构如下(示例):
+-------+----------+-------+-------+-------+-----------+---------+ |code_id|date |import1|import2|import3|postal_code|country | +-------+----------+-------+-------+-------+-----------+---------+ |1 |2024-01-01|100 |null |null |10001 |USA | |1 |2024-01-01|null |200 |null |10001 |USA | |1 |2024-01-01|null |null |300 |10001 |USA | |2 |2024-01-02|null |400 |null |SW1A 1AA |UK | +-------+----------+-------+-------+-------+-----------+---------+
- 导入Spark SQL函数库:
from pyspark.sql import functions as F
- 执行分组聚合:
result_df = data.groupBy("code_id", "date") \ .agg( # 提取每个import字段的第一个非空值 F.first("import1", ignorenulls=True).alias("import1"), F.first("import2", ignorenulls=True).alias("import2"), F.first("import3", ignorenulls=True).alias("import3"), # 保留postal_code和country,取分组内非空值(同一主键下该字段值应一致) F.first("postal_code", ignorenulls=True).alias("postal_code"), F.first("country", ignorenulls=True).alias("country") )
执行后得到的结果就是同一code_id和date下,所有import字段的非空值合并到一行,同时保留postal_code和country:
+-------+----------+-------+-------+-------+-----------+---------+ |code_id|date |import1|import2|import3|postal_code|country | +-------+----------+-------+-------+-------+-----------+---------+ |1 |2024-01-01|100 |200 |300 |10001 |USA | |2 |2024-01-02|null |400 |null |SW1A 1AA |UK | +-------+----------+-------+-------+-------+-----------+---------+
补充说明
- 如果同一分组内某个import字段存在多个非空值(比如同一主键下有两个import1的有效数据),可以用
collect_list收集所有值,再用concat_ws拼接:F.concat_ws(",", F.collect_list(F.col("import1"))).alias("import1") - 你之前用
collect_set(struct(...))效果差的原因是:这个操作会把每一行的完整结构打包成集合,最终得到一个数组类型的字段,而不是将分散的字段值合并到同一行,不符合你的需求。
内容的提问来源于stack exchange,提问作者Jersson GA
相关产品推荐
相关产品推荐

