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

Java Spark SQL:合并同Schema Dataset并实现覆盖更新

Java Spark 实现 Dataset 合并(匹配行覆盖)

针对你的需求,这里提供两种基于 Java API 的可行方案,适配 Java 11、Spark SQL 3.3.2 环境:


方案一:Left Anti Join + Union(推荐,高性能)

核心思路:先筛选出 ds1 中未在 ds2 出现过的水果行,再与 ds2 做合并,自动实现「匹配行用 ds2 覆盖,不匹配行全保留」的效果。

代码实现

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import static org.apache.spark.sql.functions.col;

// 1. 提取ds1中未在ds2出现的fruit行
Dataset<Row> ds1Unmatched = ds1.join(
    ds2,
    ds1.col("fruit").equalTo(ds2.col("fruit")),
    "left_anti"
);

// 2. 合并ds2与ds1中未匹配的行,得到最终结果
Dataset<Row> finalDs = ds2.union(ds1Unmatched);

// 可选:按fruit排序查看结果(顺序不影响业务可忽略)
finalDs.orderBy(col("fruit")).show();

方案二:Full Outer Join + Coalesce(灵活扩展)

如果后续需要处理多字段覆盖场景,可通过全外连接结合 coalesce 函数,优先取 ds2 的字段值,无匹配时 fallback 到 ds1 的值。

代码实现

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import static org.apache.spark.sql.functions.*;

// 1. 全外连接两个数据集,按fruit字段关联
Dataset<Row> joinedDs = ds1.join(
    ds2,
    ds1.col("fruit").equalTo(ds2.col("fruit")),
    "full_outer"
);

// 2. 构造最终数据集:优先选取ds2的字段值,无则取ds1的
Dataset<Row> finalDs = joinedDs.select(
    coalesce(ds2.col("fruit"), ds1.col("fruit")).alias("fruit"),
    coalesce(ds2.col("quantity"), ds1.col("quantity")).alias("quantity")
);

finalDs.show();

方案验证

针对你提供的示例数据:

  • ds2 中的 orange 会保留,ds1 中的 orange 被排除在未匹配行之外,最终 orange 的 quantity 为 5;
  • ds1 的 apple、pear、kiwi 与 ds2 的 banana、pineapple、blueberry 均会被保留,符合需求。

内容的提问来源于stack exchange,提问作者hotmeatballsoup

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 14:02:33