如何连接无公共列、行列数不同的两个Spark DataFrame
关联无公共列的Spark DataFrame解决方案
嘿,我明白你想通过添加行索引来关联这两个没有公共列、行列数还不一样的DataFrame,这个思路方向是对的,但直接用monotonically_increasing_id()可能会踩坑——因为这个函数生成的是全局唯一但不连续的64位整数,分布式环境下分区之间的ID会有间隔,导致两个DataFrame的行号没法精准对应。
给你推荐更可靠的实现方式:用row_number()窗口函数生成连续的行索引,再按需关联。
步骤1:给两个DataFrame生成连续行号
首先导入窗口函数依赖,然后给每个DF添加从1开始的连续行索引:
// 导入必要的包 import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.{row_number, lit} // 给df1添加连续行索引 val df1WithIndex = df1.withColumn( "row_index", row_number().over(Window.orderBy(lit(1))) // lit(1)表示不按实际列排序,仅生成行号 ) // 给df2添加连续行索引 val df2WithIndex = df2.withColumn( "row_index", row_number().over(Window.orderBy(lit(1))) )
注意:如果需要固定行的顺序(避免每次运行行号打乱),别用lit(1),换成你需要排序的实际列,比如df1的Customer_name或者df2的Order_ID。
步骤2:根据需求关联DataFrame
根据你想要保留的行数据,选择对应的连接类型:
- 全外连接:保留两个DF的所有行,行数不足的地方补null
val fullJoinedDF = df1WithIndex.join(df2WithIndex, Seq("row_index"), "full_outer")
- 左连接:只保留df1的所有行,df2行数不足的地方补null
val leftJoinedDF = df1WithIndex.join(df2WithIndex, Seq("row_index"), "left")
- 右连接:只保留df2的所有行,df1行数不足的地方补null
val rightJoinedDF = df1WithIndex.join(df2WithIndex, Seq("row_index"), "right")
为什么不用monotonically_increasing_id()?
这个函数是基于分区ID和分区内偏移量生成的ID,比如第一个分区生成0-99,第二个分区直接从10000开始,导致行号不连续。如果两个DF的分区数不同,它们的行索引完全没法一一对应,关联结果会完全不符合预期。
内容的提问来源于stack exchange,提问作者nirmal
相关产品推荐
相关产品推荐

