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

Spark Java API中连接DS1与DS2两个DataSet生成DS3的方法

在Spark Java API中连接两个DataSet的实现方案

没问题,我来帮你搞定Spark Java API里两个DataSet的连接操作~首先看你的DS1和DS2,它们都有Compte这个共同字段,这正是我们用来关联的核心键。下面我会一步步给你演示不同连接方式的代码实现,以及需要注意的细节。

1. 先准备强类型POJO类(可选但推荐)

Spark Java的DataSet是强类型的,定义对应DS1、DS2以及结果DS3的实体类,能让操作更清晰规范:

// 对应DS1的结构
public class CompteDetail {
    private String Compte;
    private String Lib;
    private Double ReportDebit;
    private Double ReportCredit;

    // 必须提供无参构造函数
    public CompteDetail() {}

    // 带参构造函数(方便测试模拟数据)
    public CompteDetail(String compte, String lib, Double reportDebit, Double reportCredit) {
        Compte = compte;
        Lib = lib;
        ReportDebit = reportDebit;
        ReportCredit = reportCredit;
    }

    // 所有字段的getter和setter方法
    public String getCompte() { return Compte; }
    public void setCompte(String compte) { Compte = compte; }
    public String getLib() { return Lib; }
    public void setLib(String lib) { Lib = lib; }
    public Double getReportDebit() { return ReportDebit; }
    public void setReportDebit(Double reportDebit) { ReportDebit = reportDebit; }
    public Double getReportCredit() { return ReportCredit; }
    public void setReportCredit(Double reportCredit) { ReportCredit = reportCredit; }

    // toString方法(方便打印查看结果)
    @Override
    public String toString() {
        return "CompteDetail{" +
                "Compte='" + Compte + '\'' +
                ", Lib='" + Lib + '\'' +
                ", ReportDebit=" + ReportDebit +
                ", ReportCredit=" + ReportCredit +
                '}';
    }
}

// 对应DS2的结构
public class CompteBalance {
    private String Compte;
    private Double SoldeBalance;

    public CompteBalance() {}

    public CompteBalance(String compte, Double soldeBalance) {
        Compte = compte;
        SoldeBalance = soldeBalance;
    }

    // getter和setter方法
    public String getCompte() { return Compte; }
    public void setCompte(String compte) { Compte = compte; }
    public Double getSoldeBalance() { return SoldeBalance; }
    public void setSoldeBalance(Double soldeBalance) { SoldeBalance = soldeBalance; }

    @Override
    public String toString() {
        return "CompteBalance{" +
                "Compte='" + Compte + '\'' +
                ", SoldeBalance=" + SoldeBalance +
                '}';
    }
}

// 对应连接后DS3的结构(如果需要强类型结果)
public class CompteCombined {
    private String Compte;
    private String Lib;
    private Double ReportDebit;
    private Double ReportCredit;
    private Double SoldeBalance;

    public CompteCombined() {}

    // 所有字段的getter和setter方法
    public String getCompte() { return Compte; }
    public void setCompte(String compte) { Compte = compte; }
    public String getLib() { return Lib; }
    public void setLib(String lib) { Lib = lib; }
    public Double getReportDebit() { return ReportDebit; }
    public void setReportDebit(Double reportDebit) { ReportDebit = reportDebit; }
    public Double getReportCredit() { return ReportCredit; }
    public void setReportCredit(Double reportCredit) { ReportCredit = reportCredit; }
    public Double getSoldeBalance() { return SoldeBalance; }
    public void setSoldeBalance(Double soldeBalance) { SoldeBalance = soldeBalance; }

    @Override
    public String toString() {
        return "CompteCombined{" +
                "Compte='" + Compte + '\'' +
                ", Lib='" + Lib + '\'' +
                ", ReportDebit=" + ReportDebit +
                ", ReportCredit=" + ReportCredit +
                ", SoldeBalance=" + SoldeBalance +
                '}';
    }
}

2. 核心连接操作代码

假设你已经初始化好了SparkSession,并且加载了DS1和DS2的数据源,下面是几种常见连接方式的实现:

2.1 内连接(Inner Join)

只保留**两边都有匹配Compte**的记录,这是最常用的连接方式:

// 初始化SparkSession(本地测试用,生产环境去掉master配置)
SparkSession spark = SparkSession.builder()
        .appName("JoinDataSetDemo")
        .master("local[*]")
        .getOrCreate();

// 加载DS1和DS2(这里用模拟数据举例,实际替换成你的数据源读取逻辑)
List<CompteDetail> detailList = Arrays.asList(
        new CompteDetail("447105", "Autres impôts, ta...", 77171.0, 0.0),
        new CompteDetail("753000", "Jetons de présenc...", 6839.0, 0.0),
        new CompteDetail("511107", "Valeurs à l’encai...", 0.0, 77171.0)
);
Dataset<CompteDetail> ds1 = spark.createDataset(detailList, Encoders.bean(CompteDetail.class));

List<CompteBalance> balanceList = Arrays.asList(
        new CompteBalance("447105", 992.0), // 假设完整值是992.xxx
        new CompteBalance("753000", 123.0)
);
Dataset<CompteBalance> ds2 = spark.createDataset(balanceList, Encoders.bean(CompteBalance.class));

// 执行内连接,生成Row类型的DS3
Dataset<Row> ds3Row = ds1.join(ds2, ds1.col("Compte").equalTo(ds2.col("Compte")), "inner");

// 如果需要强类型的DS3,转换为CompteCombined类型
Dataset<CompteCombined> ds3Typed = ds1.join(ds2, ds1.col("Compte").equalTo(ds2.col("Compte")))
        .select(
                ds1.col("Compte"), // 明确取ds1的Compte,避免重复字段
                ds1.col("Lib"),
                ds1.col("ReportDebit"),
                ds1.col("ReportCredit"),
                ds2.col("SoldeBalance")
        )
        .as(Encoders.bean(CompteCombined.class));

// 打印结果查看
ds3Typed.show();

2.2 左外连接(Left Outer Join)

保留DS1的所有记录,如果DS2中没有匹配的Compte,对应的SoldeBalance会填充null:

Dataset<Row> ds3LeftJoin = ds1.join(ds2, ds1.col("Compte").equalTo(ds2.col("Compte")), "left_outer");
// 同样可以转换为强类型,逻辑和内连接一致

2.3 右外连接(Right Outer Join)

保留DS2的所有记录,如果DS1中没有匹配的Compte,对应的DS1字段会填充null:

Dataset<Row> ds3RightJoin = ds1.join(ds2, ds1.col("Compte").equalTo(ds2.col("Compte")), "right_outer");

2.4 全外连接(Full Outer Join)

保留两边所有记录,没有匹配的字段会填充null:

Dataset<Row> ds3FullJoin = ds1.join(ds2, ds1.col("Compte").equalTo(ds2.col("Compte")), "full_outer");

3. 关键注意事项

  • 重复字段处理:两个DataSet都有Compte字段,连接后默认会保留两个同名字段,建议在select时明确指定取其中一个,避免后续操作报错。
  • 连接类型参数:join方法的第三个参数是连接类型字符串,可选值包括:inner、left_outer、right_outer、full_outer、cross(笛卡尔积,谨慎使用)等。
  • 空值处理:外连接会产生null值,后续可以用na().fill()方法填充默认值,比如:ds3LeftJoin.na().fill(0.0, Arrays.asList("SoldeBalance"))。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:55:56