如何在外部Hive表中覆盖指定多个分区且保留其余分区?
问题:Spark中如何批量覆盖Hive外部表的指定分区(保留其他分区)
我有一个存储在S3上的外部Hive分区表sandbox,按part字段分区,初始数据如下:
scala> q("select * from sandbox") +---+------------------------------------+----+ |id |val |part| +---+------------------------------------+----+ |0 |5d8bad52-5373-4147-8cea-9492bb0d86b9|p0 | |1 |28ed4d74-43ac-453b-bec9-651337bc18fc|p1 | |4 |3f958c0f-88a8-4afa-bcf5-f89bcead3712|p4 | |2 |cb596b60-a12a-4a19-9c37-6a12a036a71d|p2 | |3 |b69f53d6-6de4-495f-881c-9259204e4a30|p3 | +---+------------------------------------+----+
另有一张data表,存储需要更新的分区数据:
scala> q("select * from data") +---+--------------------------------------------+----+ |id |val |part| +---+--------------------------------------------+----+ |0 |updated_9521f4d0-0717-4025-b0a2-1237bf1b3b34|p0 | |1 |updated_d4987777-97f5-4676-a464-bd45877868fc|p1 | +---+--------------------------------------------+----+
需求是将data的数据合并到sandbox中,仅覆盖sandbox中与data对应的分区(p0、p1),保留其他分区(p2、p3、p4)。
尝试过以下方法,但都有问题:
- 使用
insert overwrite sandbox partition (part) select * from data:会删除sandbox中所有未被data覆盖的分区,不符合需求。 - 单分区逐个执行
insert overwrite sandbox partition(part='p1') select * from data where part = 'p1':可以实现单分区覆盖,但需要多次执行,效率低。 - 左连接合并数据集:会失去分区设计的优势,且
sandbox数据量远大于data,性能差。
完整复现过程:
scala> def q(s: String) = spark.sql(s).show(false) q: (s: String)Unit scala> case class Record(id: Int, `val`: String, part: String) defined class Record scala> q(""" | create external table sandbox ( | id int, | val string, | part string | ) using parquet | partitioned by ( | part | ) | location 's3://...../sandbox/'""") 25/03/21 19:22:28 WARN SessionState: METASTORE_FILTER_HOOK will be ignored, since hive.security.authorization.manager is set to instance of HiveAuthorizerFactory. ++ || ++ ++ scala> (0 until 10).map(i => {Record(i, java.util.UUID.randomUUID().toString(), "p" + (i % 5))}).toSeq.toDS.write.format("parquet").partitionBy("part").mode("append").saveAsTable("sandbox") scala> q(""" | create table data ( | id int, | val string, | part string | ) using parquet | partitioned by ( | part | ) | location 'hdfs:///data'""") ++ || ++ ++ scala> (0 until 2).map(i => {Record(i, "updated_" + java.util.UUID.randomUUID().toString(), "p" + (i % 5))}).toSeq.toDS.write.format("parquet").mode("append").partitionBy("part").saveAsTable("data") scala> q("select * from data") +---+--------------------------------------------+----+ |id |val |part| +---+--------------------------------------------+----+ |0 |updated_ca718ad1-319d-4bbf-9469-6539db7918f2|p0 | |1 |updated_a3e80d62-c10b-4b05-ac1e-d054f121909c|p1 | +---+--------------------------------------------+----+ scala> q("select * from sandbox") +---+------------------------------------+----+ |id |val |part| +---+------------------------------------+----+ |4 |5707218e-d717-432e-b810-a20c5037de30|p4 | |9 |5e98e6b4-474f-4674-93a0-76d26ded0679|p4 | |0 |0dd8e874-9de8-47bb-be83-59c4e6e0f021|p0 | |2 |d0ffb48e-0dca-493d-b4cb-4417d956145c|p2 | |7 |7fc72d81-7124-4ce1-94ba-da9c93233f66|p2 | |1 |53b3776c-89ce-4511-9958-9184a2a99cbf|p1 | |6 |4dcb8a74-2cb7-4b01-87e6-713d54f79b7d|p1 | |8 |312053a7-01bd-41a2-b4a6-0200a87f77b6|p3 | |5 |e6198dc2-2997-4988-99da-1fb2fa525ba7|p0 | |3 |e45489cc-bdb4-409c-972a-3229df89af73|p3 | +---+------------------------------------+----+ scala> q("insert overwrite sandbox partition (part) select * from data") ++ || ++ ++ scala> q("select * from sandbox") +---+--------------------------------------------+----+ |id |val |part| +---+--------------------------------------------+----+ |0 |updated_ca718ad1-319d-4bbf-9469-6539db7918f2|p0 | |1 |updated_a3e80d62-c10b-4b05-ac1e-d054f121909c|p1 | +---+--------------------------------------------+----+ scala> q("insert overwrite sandbox partition (part) select * from data where part = 'p0'") ++ || ++ ++ scala> q("select * from sandbox") +---+--------------------------------------------+----+ |id |val |part| +---+--------------------------------------------+----+ |0 |updated_ca718ad1-319d-4bbf-9469-6539db7918f2|p0 | +---+--------------------------------------------+----+
解决方案
方法1:启用Spark动态分区的non-strict模式 + 分区覆盖配置
这是最优方案,通过修改Spark配置,让insert overwrite仅覆盖指定分区:
- 先设置三个关键配置:
spark.sql("SET hive.exec.dynamic.partition = true") spark.sql("SET hive.exec.dynamic.partition.mode = non-strict") spark.sql("SET spark.sql.sources.partitionOverwriteMode = dynamic")
- 执行覆盖语句:
q("insert overwrite sandbox partition (part) select * from data")
此时Spark会仅覆盖data表中存在的分区(p0、p1),保留sandbox中其他未涉及的分区,完全符合需求。
配置说明:
hive.exec.dynamic.partition = true:启用动态分区,允许自动识别分区值。hive.exec.dynamic.partition.mode = non-strict:允许所有分区都是动态的,无需指定静态分区。spark.sql.sources.partitionOverwriteMode = dynamic:核心配置,将覆盖模式从"全表覆盖"改为"仅覆盖涉及的分区"。
方法2:编程批量生成单分区覆盖语句
如果无法修改集群配置,可以通过代码自动获取分区值并批量执行:
// 获取data表的所有唯一分区值 val partitions = spark.sql("select distinct part from data").collect().map(_.getString(0)) // 批量执行单分区覆盖 partitions.foreach { part => spark.sql(s""" insert overwrite sandbox partition(part='$part') select id, val, part from data where part = '$part' """) }
这种方法避免了手动逐个执行,同时保证只覆盖目标分区,适合无法修改全局配置的场景。
方法3:使用MERGE INTO(Spark 3.0+支持)
如果你的Spark版本在3.0及以上,可以用MERGE INTO实现行级的更新/插入:
spark.sql(""" MERGE INTO sandbox s USING data d ON s.part = d.part AND s.id = d.id WHEN MATCHED THEN UPDATE SET s.val = d.val WHEN NOT MATCHED THEN INSERT (id, val, part) VALUES (d.id, d.val, d.part) """)
注意:该方法是行级操作,若需要完全覆盖整个分区(而非仅更新匹配行),方法1或方法2更合适。
内容的提问来源于stack exchange,提问作者kot
相关产品推荐
相关产品推荐

