Spark Scala拆分DataFrame后合并行数异常的解决方法
解决方案:横向合并拆分后的DataFrame(保留原行数)
你的问题出在错误使用了unionByName——这个方法是纵向拼接行(把两个DF的行堆叠),所以10行+10行变成20行。而你需要的是横向拼接列,把拆分后的列合并回原结构,这时候应该用join操作,通过行标识来匹配对应行。
修正步骤:
- 给每个拆分后的DataFrame添加唯一行ID,确保每行能精准对应
- 按ID列将所有拆分DF做内连接,最后移除ID列
完整修正代码:
import org.apache.spark.sql.functions.{monotonically_increasing_id, col} // Specify the number of columns in each split val columnsPerSplit = 100 // Get the total number of columns in the original DataFrame val totalColumns = originalDF.columns.length // Calculate the number of splits required val numSplits = (totalColumns.toDouble / columnsPerSplit).ceil.toInt // Split the original DataFrame into multiple DataFrames with added row ID val splitDataFramesWithId = (0 until numSplits).map { splitIndex => val startColIndex = splitIndex * columnsPerSplit val endColIndex = Math.min((splitIndex + 1) * columnsPerSplit, totalColumns) // Select columns for the current split val selectedColumns = originalDF.columns.slice(startColIndex, endColIndex).map(col) // Add unique row ID to each split DF originalDF.select(selectedColumns: _*) .withColumn("row_id", monotonically_increasing_id()) } // Merge all split DFs by joining on row_id val mergedDF = splitDataFramesWithId.reduce((df1, df2) => df1.join(df2, Seq("row_id"), "inner") ).drop("row_id")
代码说明:
monotonically_increasing_id():生成全局唯一的行标识,确保拆分后的每行能正确匹配join(df2, Seq("row_id"), "inner"):按行ID做内连接,横向合并列,保留原行数- 最后
drop("row_id"):移除临时添加的ID列,恢复原数据结构
验证示例:
用你给出的测试案例验证:
val df1 = spark.createDataFrame(Seq((0, 1, 2)), Seq("col0", "col1", "col2")).withColumn("row_id", monotonically_increasing_id()) val df2 = spark.createDataFrame(Seq((3, 4, 5)), Seq("col3", "col4", "col5")).withColumn("row_id", monotonically_increasing_id()) df1.join(df2, Seq("row_id"), "inner").drop("row_id").show()
输出会和你的预期一致:
+----+----+----+----+----+----+ |col0|col1|col2|col3|col4|col5| +----+----+----+----+----+----+ | 0| 1| 2| 3| 4| 5| +----+----+----+----+----+----+
内容的提问来源于stack exchange,提问作者Gaurav Pargai
相关产品推荐
相关产品推荐

