Spark DataFrame对比获取匹配记录(非Spark SQL解决方案求助)
报错原因
你出现该错误的核心是:df3是仅包含empids数组列的单行DataFrame,在未和df5建立关联的情况下直接跨DataFrame引用df3的列,Spark无法在df5的元数据中找到对应字段,因此触发属性缺失的分析异常。
可行解决方案(均不依赖Spark SQL写法)
方案1:原生Spark DataFrame内连接(大数据量下性能最优)
直接使用两张表的原始DataFrame做左半连接实现存在性校验,避免将全量id收集到数组造成内存压力:
// df2为emp1表读取后的DataFrame,df5为emp2表读取后的DataFrame val matchedDf = df5.join(df2, df5("emp_no") === df2("emp_id"), "left_semi") matchedDf.show()
左半连接(left_semi)只会返回左表(df5)中在右表有匹配的记录,不会重复返回右表的字段,比普通内连接更符合你提取匹配记录的需求。
方案2:广播数组校验(适合emp1的emp_id量级较小的场景)
如果你要沿用之前收集id数组的逻辑,可以先把数组提取为本地变量再注入到df5中:
// 提取df3中的id数组为本地Scala集合 val empIdArr = df3.head().getAs[Seq[Int]]("empids") // 注入数组常量做校验 val df6 = df5.withColumn("is_exist", array_contains(lit(empIdArr), col("emp_no"))) // 过滤出存在的记录 val matchedDf = df6.filter(col("is_exist") === true).select("emp_no") matchedDf.show()
方案3:Pandas API on Spark(类Pandas语法,无Spark学习成本)
如果你更熟悉Pandas的语法,可以直接用Spark提供的Pandas API实现逻辑:
// 转为Pandas on Spark数据结构 val psEmp1 = df2.pandas_api() val psEmp2 = df5.pandas_api() // 用isin方法做存在性校验 val matchedPs = psEmp2[psEmp2["emp_no"].isin(psEmp1["emp_id"])] // 转回标准Spark DataFrame val matchedDf = matchedPs.to_spark() matchedDf.show()
执行结果
基于你提供的示例数据,最终匹配到的有效记录如下:
| emp_no |
|---|
| 1 |
| 2 |
| 3 |
内容的提问来源于stack exchange,提问作者saikiran
相关产品推荐
相关产品推荐

