Scala中Spark DataFrame基于多列的左外连接实现求助
修改后的Spark左外连接代码
要实现基于ID、date、location三列的左外连接,需要做两处关键修改:
- 确保每个参与连接的DataFrame都包含连接所需的三列,修改
getRelevantSegmentInfo函数的查询逻辑 - 更新join操作的连接条件,加入date和location列
以下是修改后的完整代码:
import org.apache.spark.sql.DataFrame import org.apache.spark.sql.SparkSession val spark = SparkSession.builder.master("local").appName("testing").enableHiveSupport().getOrCreate() import spark.implicits._ def getRelevantSegmentInfo(tableName: String, segmentName: String, pocs: Seq[String])(implicit spark: SparkSession ): DataFrame = { spark .table(tableName) .select(segmentName, "ID", "date", "location") // 修改:新增date和location列的查询 } val firstDF: DataFrame = getRelevantSegmentInfo( result_as_sequence.head.tableName, result_as_sequence.head.segmentName, result_as_sequence.head.pocs )(spark) val finalDF = result_as_sequence.tail.foldLeft(firstDF) { case (leftDF, segmentStruct) => leftDF.join( getRelevantSegmentInfo( segmentStruct.tableName, segmentStruct.segmentName, segmentStruct.pocs )(spark), Seq("ID", "date", "location"), // 修改:连接条件加入date和location "left_outer" ) }
另外修正了原代码中val conf =的错误定义,改为val spark =避免变量类型不匹配问题。
内容的提问来源于stack exchange,提问作者vicmartin
相关产品推荐
相关产品推荐

