Spark环境下Synapse湖数据库基于匹配键的数据删除覆盖问题
问题场景
- 运行环境:Azure Synapse Analytics Lake Database,Spark集群
- 数据表:
- tableA:Parquet格式,100亿行量级
- tableB:Parquet格式,1000万行量级
- 需求:删掉tableA里和tableB匹配(匹配键:Country、Year、Month、Store_cd、SKU)的1000万行,再把tableB的新数据写入tableA
- 踩坑记录:
- 直接写Spark SQL DELETE语句报错:
[UNSUPPORTED_FEATURE.TABLE_OPERATION] The feature is not supported: Table spark_catalog.dbo.table1 does not support DELETE. - 按官方建议把tableA全量加载成DataFrame再过滤,100亿行根本跑不动
- 直接写Spark SQL DELETE语句报错:
解决方案
Parquet格式本身不支持行级删除,结合两张表的量级差,给你两个高效方案:
方案一:用Spark SQL MERGE INTO(优先选,前提是表支持ACID)
如果你的Lake Database表是Delta Lake格式(Synapse Lake Database默认支持),直接用MERGE INTO能原子完成“删旧数据+插新数据”:
MERGE INTO tableA AS target USING tableB AS source ON target.Country = source.Country AND target.Year = source.Year AND target.Month = source.Month AND target.Store_cd = source.Store_cd AND target.SKU = source.SKU WHEN MATCHED THEN DELETE WHEN NOT MATCHED THEN INSERT *
Spark会自动识别tableB数据量小,把它广播到各个节点,不会全量扫描tableA,性能有保障。
方案二:广播Join+分区覆盖(适配普通Parquet表)
要是用的是普通Parquet表,不支持ACID,就按下面步骤来:
- 把tableB广播到所有Executor节点,避免重复传输小表数据
- 用left anti join筛选出tableA里不需要删的行
- 把筛选后的tableA数据和tableB合并
- 覆盖写入tableA(如果是分区表,只处理tableB涉及的分区,不用扫全量100亿行)
Scala代码示例:
import org.apache.spark.sql.functions.broadcast // 加载tableB并广播 val tableB = spark.table("tableB") val broadcastedB = broadcast(tableB) // 筛选tableA中不匹配tableB的数据 val tableAKeep = spark.table("tableA") .join(broadcastedB, Seq("Country", "Year", "Month", "Store_cd", "SKU"), "leftanti") // 合并保留数据和tableB新数据 val finalData = tableAKeep.unionByName(tableB) // 覆盖写入tableA(分区表可加partitionBy指定分区字段,比如partitionBy("Year", "Month")) finalData.write.mode("overwrite").format("parquet").saveAsTable("tableA")
核心是broadcast函数,强制Spark把小表tableB分发到每个节点,join操作本地完成,不会有大表 shuffle,性能提升明显。
内容的提问来源于stack exchange,提问作者Dan Wang
相关产品推荐
相关产品推荐

