如何通过Dataform增量模型实现BigQuery Merge时排除指定列更新
在Dataform增量模型中Merge时排除指定列(保留insert_time)
在Dataform中实现Merge时保留insert_time列的初始值,有两种直接方案,对应不同场景需求:
方案一:使用内置updateColumns配置(推荐,简洁)
Dataform增量模型支持通过updateColumns参数明确指定需要更新的列,未被列入的列(比如insert_time)在Merge的when matched阶段不会被修改,从而保留初始插入值。
示例代码
config { type: "incremental", bigquery: { partitionBy: "DATE(insert_time)" }, incremental: { uniqueKey: "id", -- 用于匹配目标表和源数据的唯一键 updateColumns: ["name", "status", "updated_at"] -- 列出所有需要更新的列,排除insert_time } } select id, name, status, current_timestamp() as updated_at, current_timestamp() as insert_time -- 仅在插入新行时生效,更新时不会覆盖目标表的insert_time from ${ref("your_source_table")} -- 增量过滤:仅获取比目标表最新数据更新的记录 where updated_at > (select coalesce(max(updated_at), timestamp('1970-01-01')) from ${self()})
关键说明
- Dataform会自动生成Merge语句,
when matched阶段仅更新updateColumns中指定的列,完全跳过insert_time。 - 新行插入时,
insert_time会被设置为当前时间,后续即使该行其他列变更,这个值也不会被覆盖。
方案二:自定义Merge语句(灵活适配复杂逻辑)
如果你更倾向于完全控制Merge的执行逻辑(和原生BigQuery Merge语法对齐),可以直接在模型中编写自定义Merge语句。
示例代码
config { type: "incremental", bigquery: { partitionBy: "DATE(insert_time)" } } -- 定义增量源数据集:仅包含需要新增或更新的记录 with incremental_source as ( select id, name, status, current_timestamp() as updated_at, current_timestamp() as insert_time from ${ref("your_source_table")} where updated_at > (select coalesce(max(updated_at), timestamp('1970-01-01')) from ${self()}) ) -- 自定义BigQuery Merge逻辑 merge into ${self()} target using incremental_source source on target.id = source.id when matched then update set target.name = source.name, target.status = source.status, target.updated_at = source.updated_at -- 刻意省略insert_time,避免更新该列 when not matched then insert (id, name, status, updated_at, insert_time) values (source.id, source.name, source.status, source.updated_at, source.insert_time)
关键说明
- 这种方式完全复刻原生BigQuery Merge的写法,你可以自由调整
when matched阶段的更新列,精准排除insert_time。 - 适合需要额外添加条件判断(比如仅当某些列发生变化时才更新)的复杂场景。
进阶:动态生成更新列(列数较多时)
如果目标表列数较多,不想手动枚举更新列,可以通过查询information_schema动态生成需要更新的列列表:
config { type: "incremental", incremental: { uniqueKey: "id", updateColumns: ${sql.concat( (select array_agg(column_name) from information_schema.columns where table_catalog = 'your_gcp_project' and table_schema = 'your_bq_dataset' and table_name = ${self().name} and column_name not in ('id', 'insert_time')) .join(", ") )} } } -- 剩余查询逻辑同方案一 select ...
注意:需要确保Dataform服务账号有权限访问BigQuery的information_schema视图。
内容的提问来源于stack exchange,提问作者alek6dj
相关产品推荐
相关产品推荐

