You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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:全量合并去重(适用于数据量不大的场景)

  1. 复制目标数据到临时存储:用Copy Activity将现有Parquet目标数据复制到临时数据集(如Blob的临时文件夹)。
  2. 追加增量数据到临时存储:用Copy Activity将步骤2提取的增量数据以「Append」模式写入同一临时数据集。
  3. 合并去重并覆盖目标:通过支持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:过滤旧数据+追加增量(适用于数据量较大的场景)

  1. 获取增量数据的主键列表:用Lookup Activity查询增量数据的所有主键,拼接成逗号分隔的字符串存入变量incremental_keys(需注意主键数据类型的格式,如字符串需加单引号)。
  2. 过滤目标数据:用Copy Activity查询现有Parquet目标数据,排除增量数据中的主键记录,写入临时数据集:
SELECT * FROM OPENROWSET(
    BULK '目标数据集路径/*.parquet',
    FORMAT = 'PARQUET'
) AS data
WHERE 你的主键列 NOT IN (@{variables('incremental_keys')})
  1. 追加增量数据:用Copy Activity将步骤2的增量数据追加到临时数据集。
  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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 10:05:39