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

如何用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       |
+-------+----------+-------+-------+-------+-----------+---------+
  1. 导入Spark SQL函数库:
from pyspark.sql import functions as F
  1. 执行分组聚合:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 13:42:50