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

Dataform增量表中重复uniqueKey行的处理逻辑与排序配置咨询

Dataform增量表相同uniqueKey多行数据的处理逻辑

默认处理逻辑

当增量更新中存在多个相同uniqueKey的行时,Dataform本身不会自动筛选保留最新行——它会直接触发对应数据仓库的合并(MERGE)操作。如果源数据里同一uniqueKey存在多行,多数数据仓库(如BigQuery、Snowflake)会抛出重复键冲突错误,因为默认的MERGE逻辑无法处理同一键的多条更新记录。

自定义合并排序逻辑的实现方式

Dataform没有直接提供“合并排序列”的配置项,但可以通过预处理源数据或自定义MERGE逻辑来实现保留最新行的需求,以下是两种常用方案:

方案1:在查询阶段先去重,保留最新行

修改SQLX文件,通过窗口函数对源数据按uniqueKey分组,仅保留last_modified最大的行,避免重复键进入合并流程:

config {
    type: "incremental",
    database: "adh-silver",
    schema: "financials",
    name: "24so_gl_transaction_actual",
    uniqueKey: ["pk_transaction_id"],
}

WITH filtered_source AS (
    SELECT
      pk_transaction_id,
      transaction_number,
      transaction_date,
      financial_period,
      amount,
      last_modified,
      created_time,
      -- 按主键分组,标记最新行
      ROW_NUMBER() OVER (PARTITION BY pk_transaction_id ORDER BY last_modified DESC) AS rn
    FROM
      ${ref("adh-silver", "financials", "24so_gl_transaction_historic")}
    ${(
        when(
            incremental(),
            `WHERE last_modified > (SELECT MAX(last_modified) FROM ${self()})`
        )
    )}
)
SELECT
  pk_transaction_id,
  transaction_number,
  transaction_date,
  financial_period,
  amount,
  last_modified,
  created_time
FROM filtered_source
WHERE rn = 1 -- 仅保留每个主键对应的最新行

方案2:自定义MERGE逻辑(进阶)

如果需要更灵活的冲突处理规则(比如仅当源行更新时间晚于目标行时才更新),可以禁用Dataform默认的增量合并逻辑,手动编写MERGE语句:

config {
    type: "incremental",
    database: "adh-silver",
    schema: "financials",
    name: "24so_gl_transaction_actual",
    uniqueKey: ["pk_transaction_id"],
    incrementalMerge: false -- 禁用默认合并逻辑
}

-- 增量更新时执行自定义MERGE
${when(incremental(), `
MERGE INTO ${self()} AS target
USING (
    SELECT
      pk_transaction_id,
      transaction_number,
      transaction_date,
      financial_period,
      amount,
      last_modified,
      created_time
    FROM ${ref("adh-silver", "financials", "24so_gl_transaction_historic")}
    WHERE last_modified > (SELECT MAX(last_modified) FROM ${self()})
) AS source
ON target.pk_transaction_id = source.pk_transaction_id
WHEN MATCHED AND source.last_modified > target.last_modified THEN
  UPDATE SET
    transaction_number = source.transaction_number,
    transaction_date = source.transaction_date,
    financial_period = source.financial_period,
    amount = source.amount,
    last_modified = source.last_modified,
    created_time = source.created_time
WHEN NOT MATCHED THEN
  INSERT (pk_transaction_id, transaction_number, transaction_date, financial_period, amount, last_modified, created_time)
  VALUES (source.pk_transaction_id, source.transaction_number, source.transaction_date, source.financial_period, source.amount, source.last_modified, source.created_time);
`)}

-- 全量刷新时的逻辑
${when(not(incremental()), `
SELECT
  pk_transaction_id,
  transaction_number,
  transaction_date,
  financial_period,
  amount,
  last_modified,
  created_time
FROM ${ref("adh-silver", "financials", "24so_gl_transaction_historic")}
`)}

内容的提问来源于stack exchange,提问作者Chriklev

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 07:04:51