Spark Scala中嵌套列Join结果不符合预期问题排查
Spark Scala中嵌套列自连接结果不符合预期的问题分析
问题现象
在Spark Scala中对Dataset的嵌套列执行自连接时,返回了不符合连接条件的笛卡尔积结果,而非预期的匹配行。
示例代码及结果
case class A(id: String) case class B(randomField: String) val ds1 = spark.createDataset(Seq(A("5"), A("4"))) val ds2 = spark.createDataset(Seq(B("4"), B("5"))) val ds3 = ds1.joinWith(ds2, ds1("id") === ds2("randomField")) ds3.show /* +---+---+ | _1| _2| +---+---+ |{5}|{5}| |{4}|{4}| +---+---+ */ // 错误的自连接写法 ds3 .joinWith( ds3, ds3("_1.id") === ds3("_1.id"), "inner" ) /* 实际返回结果(笛卡尔积): +----------+----------+ | _1| _2| +----------+----------+ |{{5}, {5}}|{{5}, {5}}| |{{5}, {5}}|{{4}, {4}}| |{{4}, {4}}|{{5}, {5}}| |{{4}, {4}}|{{4}, {4}}| +----------+----------+ 预期结果: +----------+----------+ | _1| _2| +----------+----------+ |{{5}, {5}}|{{5}, {5}}| |{{4}, {4}}|{{4}, {4}}| +----------+----------+ */
问题原因
你在自连接时直接使用ds3("_1.id") === ds3("_1.id")作为连接条件,Spark无法区分左右两个逻辑数据集的列——这个条件会被解析为当前行的_1.id等于自身,也就是永远为true,因此会触发笛卡尔积(2行×2行=4行),而非按id匹配的连接。
验证:提取嵌套列到根级别可正常执行
当把嵌套列提取到根级别并分别命名后,Spark能明确区分左右表的列,连接条件生效:
ds3.withColumn("id1", col("_1.id")) .joinWith( ds3.withColumn("id2", col("_1.id")), col("id1") === col("id2"), "inner" ) /* 返回结果: +-------------+-------------+ | _1| _2| +-------------+-------------+ |{{5}, {5}, 5}|{{5}, {5}, 5}| |{{4}, {4}, 4}|{{4}, {4}, 4}| +-------------+-------------+ */
正确的自连接写法
给左右两个数据集分别设置别名,通过别名明确引用嵌套列,让Spark能区分左右表的字段:
ds3.alias("left") .joinWith( ds3.alias("right"), col("left._1.id") === col("right._1.id"), "inner" )
执行后会得到你预期的匹配结果。
内容的提问来源于stack exchange,提问作者849856487764
相关产品推荐
相关产品推荐

