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
相关产品推荐
相关产品推荐

