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

