Azure Synapse无服务器SQL能否合并Parquet文件实现Upsert?
解决方案
一、合并全量与增量Parquet文件并去重,重建外部表
Synapse Serverless的外部表只是存储文件的元数据映射,本身无法执行DML操作,所以要实现去重合并,得先处理存储上的文件,再重建外部表,以下三种方式都可行:
1. 用Synapse Spark池处理
这是数据量较大场景下的灵活方案:
- 将全量和增量Parquet文件分别读入Spark DataFrame
- 基于主键去重,保留增量里的最新数据(比如按
SYS_CHANGE_VERSION排序后去重,增量数据会覆盖全量的重复项) - 将合并后的DataFrame写入存储Gen2的新路径(可直接覆盖旧全量路径,或写入新路径)
- 最后用CETAS重新创建外部表指向这个新路径
示例Scala代码:
// 读取全量和增量数据 val full_df = spark.read.parquet("abfss://container@yourstorage.dfs.core.windows.net/full_data/") val incr_df = spark.read.parquet("abfss://container@yourstorage.dfs.core.windows.net/incr_data/") // 合并后按主键去重,保留版本号最大的记录 val combined_df = full_df.union(incr_df) .orderBy($"SYS_CHANGE_VERSION".desc) .dropDuplicates(Seq("your_primary_key")) // 写入目标路径(overwrite模式会清空旧文件) combined_df.write.mode("overwrite").parquet("abfss://container@yourstorage.dfs.core.windows.net/updated_full_data/")
2. 用Serverless SQL的CETAS结合逻辑去重
直接通过SQL完成合并去重,无需额外计算资源:
- 先创建临时外部表分别映射全量和增量文件
- 用窗口函数按主键分组,取
SYS_CHANGE_VERSION最大的记录,再通过CETAS写入新的存储路径 - 最终用这个新路径的外部表作为更新后的ODS表
示例SQL:
-- 定义存储数据源和Parquet格式(已定义可跳过) CREATE EXTERNAL DATA SOURCE storage_gen2 WITH ( LOCATION = 'abfss://container@yourstorage.dfs.core.windows.net/' ); CREATE EXTERNAL FILE FORMAT parquet_format WITH ( FORMAT_TYPE = PARQUET, DATA_COMPRESSION = 'org.apache.hadoop.io.compress.SnappyCodec' ); -- 创建临时外部表映射全量、增量文件 CREATE EXTERNAL TABLE #full_temp WITH (LOCATION = 'full_data/', DATA_SOURCE = storage_gen2, FILE_FORMAT = parquet_format) AS SELECT * FROM OPENROWSET(BULK 'full_data/*.parquet', FORMAT = 'PARQUET') AS f; CREATE EXTERNAL TABLE #incr_temp WITH (LOCATION = 'incr_data/', DATA_SOURCE = storage_gen2, FILE_FORMAT = parquet_format) AS SELECT * FROM OPENROWSET(BULK 'incr_data/*.parquet', FORMAT = 'PARQUET') AS i; -- CETAS写入去重后的合并数据到新路径 CREATE EXTERNAL TABLE ods_updated WITH (LOCATION = 'updated_full_data/', DATA_SOURCE = storage_gen2, FILE_FORMAT = parquet_format) AS SELECT * FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY your_primary_key ORDER BY SYS_CHANGE_VERSION DESC) AS rn FROM ( SELECT * FROM #full_temp UNION ALL SELECT * FROM #incr_temp ) AS all_records ) AS ranked_records WHERE rn = 1;
3. 用ADF数据流处理
如果已经在用ADF做复制活动,可直接用数据流完成合并:
- 新增两个数据源,分别指向全量和增量Parquet文件
- 用
Union组件合并两个数据集 - 用
Window组件按主键分组,按SYS_CHANGE_VERSION降序生成行号 - 用
Filter组件只保留行号=1的记录(即最新数据) - 将结果写入存储Gen2的目标路径(覆盖旧全量路径)
- 最后在ADF管道中添加存储过程活动,执行CETAS重建外部表(若路径未变,Serverless会自动读取新文件,无需重建)
二、解决复制活动Upsert报错的问题
报错核心原因:Serverless外部表是只读的,不支持INSERT/UPSERT操作,因此不能把外部表作为复制活动的目标。调整方案如下:
- 将复制活动的目标改为存储Gen2的路径(比如临时增量路径
incr_temp/),不要指向外部表 - 用上述任意一种合并方式,把临时路径的增量数据和正式全量路径的数据合并去重,写入正式路径
- 若外部表的LOCATION未变化,Serverless会自动读取新文件;若路径变更,用CETAS重建外部表
另外,若想简化流程,可直接用ADF数据流对存储上的全量Parquet文件做Upsert:数据流支持直接对Parquet文件按主键匹配,完成更新或插入操作,无需单独保留增量文件,更新后外部表自动映射最新数据。
内容的提问来源于stack exchange,提问作者Jean-Christophe Rat-Patron
相关产品推荐
相关产品推荐

