如何利用Copy Activity预复制脚本基于源Change Tracking表删除Sink端记录?
解决批量表场景下Change Tracking增量复制的Upsert/删除问题
我之前在处理从SQL Server到Azure SQL的批量表增量复制时,也碰到过和你一模一样的痛点——Copy Activity默认的追加模式没法处理更新和删除,挨个写存储过程又太繁琐。分享几个不用手动为每个表造存储过程的方案:
方案1:利用变更跟踪元数据写通用预复制脚本
你可以通过SQL Server的变更跟踪系统视图,动态生成针对目标库的删除语句,在Copy Activity执行前先清理掉需要更新或删除的记录。
核心思路是:
- 先获取本次同步的起始变更版本(比如从你存储的上次同步版本号)
- 遍历所有开启了变更跟踪的表,从
CHANGETABLE(CHANGES)中提取出操作类型为U(更新)或D(删除)的主键值 - 动态拼接DELETE语句,在目标Azure SQL库中删除对应主键的记录
示例动态SQL脚本(可以在预复制脚本中执行,注意替换变量):
DECLARE @last_sync_version BIGINT = 123; -- 替换为你存储的上次同步版本 DECLARE @target_db NVARCHAR(128) = 'YourAzureSQLDB'; -- 目标数据库名 DECLARE @table_cursor CURSOR; DECLARE @table_name NVARCHAR(128), @schema_name NVARCHAR(128), @primary_key NVARCHAR(128); -- 遍历所有开启变更跟踪的表 SET @table_cursor = CURSOR FOR SELECT t.name AS table_name, s.name AS schema_name, c.name AS primary_key_column FROM sys.tables t JOIN sys.schemas s ON t.schema_id = s.schema_id JOIN sys.change_tracking_tables ctt ON t.object_id = ctt.object_id JOIN sys.key_constraints kc ON t.object_id = kc.parent_object_id JOIN sys.columns c ON kc.parent_object_id = c.object_id AND c.column_id = kc.key_columns[1] WHERE kc.type = 'PK'; OPEN @table_cursor; FETCH NEXT FROM @table_cursor INTO @table_name, @schema_name, @primary_key; WHILE @@FETCH_STATUS = 0 BEGIN DECLARE @delete_sql NVARCHAR(MAX); -- 动态生成删除目标表中对应变更记录的SQL SET @delete_sql = N' USE ' + QUOTENAME(@target_db) + N'; DELETE FROM ' + QUOTENAME(@schema_name) + N'.' + QUOTENAME(@table_name) + N' WHERE ' + QUOTENAME(@primary_key) + N' IN ( SELECT ' + QUOTENAME(@primary_key) + N' FROM CHANGETABLE(CHANGES ' + QUOTENAME(@schema_name) + N'.' + QUOTENAME(@table_name) + N', ' + CAST(@last_sync_version AS NVARCHAR) + N') ct WHERE ct.SYS_CHANGE_OPERATION IN (''U'', ''D'') );'; EXEC sp_executesql @delete_sql; FETCH NEXT FROM @table_cursor INTO @table_name, @schema_name, @primary_key; END CLOSE @table_cursor; DEALLOCATE @table_cursor;
注意:这个脚本假设每个表只有单个主键列,如果是复合主键,需要调整主键列的提取逻辑。另外,要确保执行脚本的账号在源库和目标库都有足够权限。
方案2:使用ADF Copy Activity的Upsert模式(推荐)
如果你的Azure Data Factory是v2及以上版本,其实可以直接利用Copy Activity的Upsert写入行为,不用自己写预复制脚本或存储过程:
- 在Copy Activity的Sink设置中,将
Write behavior改为Upsert - 指定目标表的主键列(和源表主键一致)
- 源端依然用Change Tracking查询获取所有变更记录(包括I/U/D)
ADF会自动将源端的变更记录和目标表进行匹配:
- 对于
I(插入)操作:直接追加 - 对于
U(更新)操作:更新目标表中对应主键的记录 - 对于
D(删除)操作:需要额外配置——因为Upsert默认不处理删除,你可以在Copy Activity之后加一个Lookup Activity获取变更中的删除记录,再用Stored Procedure Activity执行删除,或者结合方案1的动态脚本处理删除。
这个方案配置简单,不用写大量SQL,适合大多数批量表场景。
方案3:动态生成通用合并存储过程
如果还是倾向于用存储过程的方式,可以写一个动态SQL脚本,自动为所有开启变更跟踪的表生成对应的合并存储过程和表类型,不用手动逐个创建:
- 遍历系统视图获取表结构、主键信息
- 动态创建表类型(用于传递变更记录)
- 动态生成MERGE语句的存储过程,处理I/U/D操作
示例生成脚本的核心片段:
DECLARE @table_name NVARCHAR(128), @schema_name NVARCHAR(128), @pk_columns NVARCHAR(MAX); -- 遍历表获取主键等信息(省略游标逻辑,类似方案1) -- 动态创建表类型 SET @create_type_sql = N' CREATE TYPE ' + QUOTENAME(@schema_name) + N'.' + QUOTENAME(@table_name + '_ChangeType') + N' AS TABLE ( ' + @column_definitions + N', SYS_CHANGE_OPERATION NCHAR(1) );'; EXEC sp_executesql @create_type_sql; -- 动态生成合并存储过程 SET @create_proc_sql = N' CREATE PROCEDURE ' + QUOTENAME(@schema_name) + N'.' + QUOTENAME('Merge_' + @table_name) + N' @changes ' + QUOTENAME(@schema_name) + N'.' + QUOTENAME(@table_name + '_ChangeType') + N' READONLY AS BEGIN MERGE INTO ' + QUOTENAME(@schema_name) + N'.' + QUOTENAME(@table_name) + N' t USING @changes s ON t.' + @pk_columns + N' = s.' + @pk_columns + N' WHEN MATCHED AND s.SYS_CHANGE_OPERATION = ''U'' THEN UPDATE SET ' + @update_set_clause + N' WHEN NOT MATCHED AND s.SYS_CHANGE_OPERATION = ''I'' THEN INSERT (' + @insert_columns + N') VALUES (' + @insert_values + N') WHEN MATCHED AND s.SYS_CHANGE_OPERATION = ''D'' THEN DELETE; END;'; EXEC sp_executesql @create_proc_sql;
这个脚本一次性生成所有需要的存储过程和表类型,后续新增表时再跑一次即可。
总结
- 不需要为每个表手动创建存储过程和表类型,用动态SQL或ADF的Upsert功能就能搞定
- 优先推荐方案2(ADF Upsert),配置最简便
- 如果必须用预复制脚本,方案1的动态SQL可以批量处理所有表的删除/更新前置操作
内容的提问来源于stack exchange,提问作者Werner Jongenelis
相关产品推荐
相关产品推荐

