Spark如何仅覆盖存在更新的指定分区来更新历史数据
Spark分区表增量更新仅重写变更分区解决方案
核心逻辑
需求为按deptno分区的历史表仅重写增量数据涉及的分区,避免全表重写,核心实现逻辑为:
- 定位增量数据涉及的所有
deptno分区 - 仅读取历史表中这些待更新分区的数据
- 对变更分区数据做新旧合并
- 开启Spark动态分区覆盖配置,仅重写变更分区
实现代码
import org.apache.spark.sql.functions.coalesce // 开启动态分区覆盖核心配置 spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic") // 读取增量更新数据 val finalQry = "SELECT * FROM emp" val df2 = spark.read.format("jdbc") .option("url", "jdbc:mysql://localhost:3306/test") .option("user", "root") .option("dbtable", "(" + finalQry + ") as t") .load() // 提取增量数据涉及的所有需要更新的分区 val updateDepts = df2.select("deptno").distinct().as[Int].collect() // 仅读取历史表中待更新分区的数据,避免全表扫描 val df1 = spark.read.orc("data/emp/").filter(col("deptno").isin(updateDepts:_*)) // 合并新旧数据:empno匹配优先用增量数据,否则用历史数据 val projections = df1.schema.fields.map { field => coalesce(df2.col(field.name), df1.col(field.name)).as(field.name) } val updatedPartitionDf = df1.join(df2, df1.col("empno") === df2.col("empno"), "fullouter") .select(projections: _*) // 写入时仅覆盖变更分区,原有其他分区完全保留 updatedPartitionDf.write .format("orc") .mode("overwrite") .partitionBy("deptno") .save("data/emp/")
关键说明
- 动态分区覆盖配置
spark.sql.sources.partitionOverwriteMode必须设置为dynamic,默认值为static,静态模式下会覆盖全部分区 - 仅读取待更新分区的历史数据,大幅降低历史表的IO开销,适合海量历史表场景
- 该方案同时兼容增量数据包含新增分区、新增员工记录的场景
内容的提问来源于stack exchange,提问作者Aravind Kumar Anugula
相关产品推荐
相关产品推荐

