Spark-SQL连接具有相同列名的两个数据框/数据集

New*_*ies 5 java apache-spark apache-spark-sql apache-spark-dataset

我有以下两个数据集

controlSetDF : has columns loan_id, merchant_id, loan_type, created_date, as_of_date
accountDF : has columns merchant_id, id, name, status, merchant_risk_status
Run Code Online (Sandbox Code Playgroud)

我正在使用Java Spark API来加入它们,我只需要最终数据集中的特定列

private String[] control_set_columns = {"loan_id", "merchant_id", "loan_type"};
private String[] sf_account_columns = {"id as account_id", "name as account_name", "merchant_risk_status"};

controlSetDF.selectExpr(control_set_columns)                                               
.join(accountDF.selectExpr(sf_account_columns),controlSetDF.col("merchant_id").equalTo(accountDF.col("merchant_id")), 
"left_outer"); 
Run Code Online (Sandbox Code Playgroud)

但我得到以下错误

org.apache.spark.sql.AnalysisException: resolved attribute(s) merchant_id#3L missing from account_name#131,loan_type#105,account_id#130,merchant_id#104L,loan_id#103,merchant_risk_status#2 in operator !Join LeftOuter, (merchant_id#104L = merchant_id#3L);;!Join LeftOuter, (merchant_id#104L = merchant_id#3L)
Run Code Online (Sandbox Code Playgroud)

似乎存在问题,因为两个数据帧都具有merchant_id列。

注意:如果我不使用.selectExpr(),它将正常工作。但是它将显示第一和第二数据集的所有列。

Sil*_*vio 2

如果两个 DataFrame 中的连接列名称相同,您可以简单地将其定义为连接条件。在 Scala 中,它更清晰一些,而在 Java 中,您需要将 Java List 转换为 Scala Seq:

Seq<String> joinColumns = scala.collection.JavaConversions
  .asScalaBuffer(Lists.newArrayList("merchant_id"));

controlSetDF.selectExpr(control_set_columns)
  .join(accountDF.selectExpr(sf_account_columns), joinColumns), "left_outer");
Run Code Online (Sandbox Code Playgroud)

这将导致 DataFrame 仅包含一个连接列。