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
相关产品推荐
相关产品推荐

