Spark Java环境下如何合并Schema不同的两个Parquet文件
解决方法
你需要的不是简单的行拼接,而是按epochMillis作为主键关联两个数据集,同主键时优先取第二个数据集的重叠列值,缺失列补null,所以应该用全外连接+列合并的逻辑实现,具体Java代码如下:
import org.apache.spark.sql.functions; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; // 读取两个Parquet文件,分别设置别名 Dataset<Row> df1 = testSparkSession.read().option("mergeSchema",true).parquet("D:\\ABC\\abc.parquet").alias("df1"); Dataset<Row> df2 = testSparkSession.read().option("mergeSchema",true).parquet("D:\\EFG\\efg.parquet").alias("df2"); // 按epochMillis字段进行全外连接,保留两个表的所有行 Dataset<Row> joinedDf = df1.join(df2, df1.col("epochMillis").equalTo(df2.col("epochMillis")), "full_outer"); // 按需求规则合并列,得到最终结果 Dataset<Row> resultDf = joinedDf.select( // 主键列取两个表中不为空的任意值 functions.coalesce(df2.col("epochMillis"), df1.col("epochMillis")).alias("epochMillis"), // 公共列优先取df2的值,df2无值时取df1的值 functions.coalesce(df2.col("one"), df1.col("one")).alias("one"), functions.coalesce(df2.col("two"), df1.col("two")).alias("two"), functions.coalesce(df2.col("three"), df1.col("three")).alias("three"), // 保留两个表各自的独有列,无值时自动为null df1.col("four").alias("four"), df2.col("five").alias("five") ); // 可选:按epochMillis排序,和示例输出顺序一致 resultDf = resultDf.orderBy("epochMillis"); // 验证结果 resultDf.show();
逻辑说明
- 全外连接
full_outer可以保留两个数据集所有的epochMillis记录,不会出现数据丢失 coalesce函数会按参数顺序返回第一个非空值,刚好实现同主键时优先取第二个数据集的重叠列值的需求- 两个数据集各自独有的列会自动保留,对应另一个数据集不存在的行自动填充null,完全匹配你给出的期望结果
之前方法报错原因
union()要求两个数据集的列数、列顺序完全一致,你的两个数据集列数不同所以报错- 低版本Spark的
unionByName()默认不允许存在仅在单个数据集里的列,就算你用Spark 2.3+版本开启allowMissingColumns=true参数,unionByName()实现的是行拼接,不会按主键合并行,也不符合你的需求。
内容的提问来源于stack exchange,提问作者Developer208
相关产品推荐
相关产品推荐

