Spark Scala中两个DataFrame非匹配记录的左外连接实现
在Spark Scala中筛选左外连接的非匹配记录
首先,我们先明确你的两个DataFrame数据:
DataFrame 1(左表)
+-------------+-------------------------+--------------+--------+----------+-----------------------+---------------------+-------------------+-----------------------+--------------------------+--------------------------+-----------+ |DataPartition|TimeStamp |OrganizationID|SourceID|_auditorId|sr:AuditorEnumerationId|sr:AuditorOpinionCode|sr:AuditorOpinionId|sr:IsPlayingAuditorRole|sr:IsPlayingCSRAuditorRole|sr:IsPlayingTaxAdvisorRole|FFAction | +-------------+-------------------------+--------------+--------+----------+-----------------------+---------------------+-------------------+-----------------------+--------------------------+--------------------------+-----------+ |Japan |2018-05-03T09:52:48+00:00|4295876589 |195 |null |null |null |null |null |null |null |O | |Japan |2018-05-03T08:10:19+00:00|4295876589 |196 |null |null |null |null |null |null |null |D | |Japan |2018-05-03T09:52:48+00:00|4295876589 |194 |null |null |null |null |null |null |null |O | +-------------+-------------------------+--------------+--------+----------+-----------------------+---------------------+-------------------+-----------------------+--------------------------+--------------------------+-----------+
DataFrame 2(右表)
DataPartition TimeStamp OrganizationID SourceID _auditorId sr:AuditorEnumerationId sr:AuditorOpinionCode sr:AuditorOpinionId sr:IsPlayingAuditorRole sr:IsPlayingCSRAuditorRole sr:IsPlayingTaxAdvisorRole FFAction Japan 2018-05-03T08:06:06+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T08:06:06+00:00 4295876589 195 16157 1002485247 UWE 3010547 true false false O Japan 2018-05-03T09:48:33+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T09:48:33+00:00 4295876589 195 16157 1002485247 UWE 3010547 true false false O Japan 2018-05-03T07:27:10+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T07:27:10+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-03T07:27:10+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O Japan 2018-05-03T07:35:42+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T07:35:42+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-03T07:35:42+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O Japan 2018-05-03T09:34:46+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T09:34:46+00:00 4295876589 195 16157 1002485247 UWE 3010547 true false false O Japan 2018-05-03T08:10:19+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T08:10:19+00:00 4295876589 195 16157 1002485247 UWE 3010547 true false false O Japan 2018-05-03T07:28:16+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T07:28:16+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-03T07:28:16+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O Japan 2018-05-02T09:05:04+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-02T09:05:04+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-02T09:05:04+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O Japan 2018-05-03T07:31:28+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T07:31:28+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-03T07:31:28+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O Japan 2018-05-03T07:22:58+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T07:22:58+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-03T07:22:58+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O Japan 2018-05-03T09:45:22+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T09:45:22+00:00 4295876589 195 16157 1002485247 UWE 3010547 true false false O Japan 2018-05-03T07:11:26+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T07:11:26+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-03T07:11:26+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O Japan 2018-05-03T07:00:45+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T07:00:45+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-03T07:00:45+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O Japan 2018-05-03T07:36:47+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O Japan 2018-05-03T07:36:47+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O Japan 2018-05-03T07:36:47+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true
解决方案代码
要实现左外连接并筛选非匹配记录,我们可以按以下步骤操作:
import org.apache.spark.sql.functions.col // 假设你的两个DataFrame分别命名为df1(左表)和df2(右表) val joinKeys = Seq("DataPartition", "OrganizationID", "SourceID", "TimeStamp") // 执行左外连接:保留df1的所有行,匹配df2的行,不匹配的df2字段为null val leftJoinedDF = df1.join(df2, joinKeys, "left_outer") // 筛选左表中没有匹配右表的记录:通过检查右表的非空字段(比如_auditorId)是否为null val nonMatchingRecords = leftJoinedDF.filter(col("_auditorId").isNull) // 如果你只需要保留左表的原始字段,可以选择如下方式 val nonMatchingDF1Only = nonMatchingRecords.select(df1.columns.map(col): _*) // 查看结果 nonMatchingDF1Only.show()
代码解释
- 连接键选择:我们使用
DataPartition、OrganizationID、SourceID和TimeStamp作为连接键,因为这些字段的组合能唯一标识每条记录,确保我们匹配的是完全一致的条目。 - 左外连接逻辑:左外连接会保留左表(df1)的所有行,右表(df2)中匹配的行会填充对应字段,不匹配的行则右表所有字段为
null。 - 筛选非匹配记录:右表中的
_auditorId字段没有null值(从你的数据可以看到),所以当这个字段为null时,说明左表的该行在右表中没有匹配项,这正是我们需要筛选的记录。
运行这段代码后,你会得到df1中所有在df2里找不到匹配的行,也就是你要的非匹配记录。
内容的提问来源于stack exchange,提问作者Atharv Thakur
相关产品推荐
相关产品推荐

