如何通过Copy Activity实现从Oracle到Parquet的SCD1增量加载?
用Copy Activity实现Oracle到Parquet的SCD1增量加载(基于水印列
Rma_date) SCD1的核心是保留最新数据,增量加载需同步上次同步后新增/更新的记录,并替换目标中的旧数据。由于Parquet是文件存储不支持行级Upsert,结合Copy Activity的实现流程如下:
步骤1:维护同步状态(记录上次同步的水印值)
- 准备一个同步状态存储:可以用小型数据库表(如Azure SQL DB)或Blob存储的文本文件,用于记录每次同步完成后的最大
Rma_date。首次同步前,将初始值设为极小值(如'1900-01-01')。 - 使用Lookup Activity读取该状态值,存入变量
last_sync_dt,供后续增量过滤使用。
步骤2:提取Oracle增量数据
- 在Copy Activity的源配置中,使用带过滤条件的SQL查询获取增量数据,替换
YOUR_ORACLE_TABLE为实际表名,确保时间格式与Oracle的Rma_date类型匹配:
SELECT * FROM YOUR_ORACLE_TABLE WHERE Rma_date > TO_TIMESTAMP('@{variables('last_sync_dt')}', 'YYYY-MM-DD HH24:MI:SS')
- 同时,用另一个Lookup Activity获取本次同步的最大水印值,用于更新同步状态:
SELECT NVL(MAX(Rma_date), TO_TIMESTAMP('@{variables('last_sync_dt')}', 'YYYY-MM-DD HH24:MI:SS')) AS current_max_dt FROM YOUR_ORACLE_TABLE WHERE Rma_date > TO_TIMESTAMP('@{variables('last_sync_dt')}', 'YYYY-MM-DD HH24:MI:SS')
将结果存入变量current_max_dt,若本次无增量数据,该值与last_sync_dt一致。
步骤3:实现Parquet目标的SCD1更新(核心)
由于Parquet不支持行级更新,需通过「合并去重+覆盖写入」实现:
方案A:全量合并去重(适用于数据量不大的场景)
- 复制目标数据到临时存储:用Copy Activity将现有Parquet目标数据复制到临时数据集(如Blob的临时文件夹)。
- 追加增量数据到临时存储:用Copy Activity将步骤2提取的增量数据以「Append」模式写入同一临时数据集。
- 合并去重并覆盖目标:通过支持Parquet查询的服务(如Azure Synapse Serverless)执行去重SQL,保留每个主键对应的最新记录,再用Copy Activity将结果写入最终目标(覆盖模式):
SELECT t.* FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY 你的主键列 ORDER BY Rma_date DESC) AS rn FROM OPENROWSET( BULK '临时数据集路径/*.parquet', FORMAT = 'PARQUET' ) AS data ) t WHERE t.rn = 1
替换你的主键列为实际表的主键字段(如id)。
方案B:过滤旧数据+追加增量(适用于数据量较大的场景)
- 获取增量数据的主键列表:用Lookup Activity查询增量数据的所有主键,拼接成逗号分隔的字符串存入变量
incremental_keys(需注意主键数据类型的格式,如字符串需加单引号)。 - 过滤目标数据:用Copy Activity查询现有Parquet目标数据,排除增量数据中的主键记录,写入临时数据集:
SELECT * FROM OPENROWSET( BULK '目标数据集路径/*.parquet', FORMAT = 'PARQUET' ) AS data WHERE 你的主键列 NOT IN (@{variables('incremental_keys')})
- 追加增量数据:用Copy Activity将步骤2的增量数据追加到临时数据集。
- 覆盖写入目标:用Copy Activity将临时数据集的内容覆盖写入最终Parquet目标。
步骤4:更新同步状态
用Copy Activity或Stored Procedure Activity将变量current_max_dt的值写入同步状态存储,确保下次同步能基于最新的水印值过滤增量数据。
关键注意事项
- 时间格式一致性:Oracle的
Rma_date若为TIMESTAMP类型,需确保变量格式与TO_TIMESTAMP的参数完全匹配,避免转换错误。 - 主键准确性:必须明确表的主键字段,否则无法正确去重并保留最新数据。
- 空增量处理:若步骤2未获取到增量数据,可跳过步骤3的合并操作,直接更新同步状态即可。
内容的提问来源于stack exchange,提问作者AzSurya Teja
相关产品推荐
相关产品推荐

