如何在Spark多分区层级下使用SaveMode.Ignore写入数据?
问题解答
1. 原生SaveMode.Ignore不支持按全部分区层级检查
Spark自带的SaveMode.Ignore逻辑非常简单:只要目标根目录outputPath已经存在,就直接跳过整个写入操作,根本不会去校验具体的id/dt/tm分区是否已存在。哪怕你本次要写入的分区全是全新的,只要根目录有数据,就完全不执行写入,这显然不符合你的需求。
2. Spark原生可用的替代方案
Spark没有直接提供按分区维度跳过已存在数据的内置模式,但可以通过两种原生方案实现类似效果:
方案一:过滤已存在分区后写入
先加载目标路径下已有的分区元数据,从rawData中过滤掉已经存在的(id, dt, tm)组合,再执行写入。示例代码:
// 读取已存在的分区列并去重 val existingPartitions = spark.read.parquet(outputPath) .select("id", "dt", "tm") .distinct() // 只保留rawData中不存在的分区数据 val dataToWrite = rawData.join(existingPartitions, Seq("id", "dt", "tm"), "left_anti") // 用Append模式写入过滤后的数据 dataToWrite.coalesce(5000) .write .partitionBy("id", "dt", "tm") .mode(SaveMode.Append) .parquet(outputPath)
如果目标路径数据量很大,读取全部分区元数据性能开销高,可以直接用Hadoop的FileSystem API遍历S3上的分区目录,解析出已存在的(id, dt, tm)集合,效率会更高。
方案二:动态分区覆盖
开启spark.sql.sources.partitionOverwriteMode=dynamic配置后,使用SaveMode.Overwrite写入时,只会覆盖本次写入包含的分区,不会影响其他分区。虽然是覆盖而非跳过,但能保证ETL的幂等性,重复执行也不会破坏已有数据。示例代码:
// 开启动态分区覆盖配置 spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic") rawData.coalesce(5000) .write .partitionBy("id", "dt", "tm") .mode(SaveMode.Overwrite) .parquet(outputPath)
这个方案不需要提前读取已有分区,性能更优,但如果你的需求是完全不触碰已存在的分区数据,这个方案就不适用。
3. 自行实现的思路
如果上述方案都不符合要求,可以自己实现分区检查逻辑:
- 用Hadoop的
FileSystemAPI遍历outputPath下的id/dt/tm层级目录; - 将目录路径解析为
(id, dt, tm)的集合; - 在Spark中过滤
rawData,只保留不在该集合中的分区; - 最后用
SaveMode.Append写入过滤后的数据。
这种方式直接操作文件系统目录,不需要加载数据到Spark内存,性能比读取全部分区元数据更好。
内容的提问来源于stack exchange,提问作者nirkov
相关产品推荐
相关产品推荐

