基于DBT BigQuery的insert_overwrite增量逻辑实现求助
按文件名实现增量数据替换逻辑解决方案
需求说明
我们获取的文件名称带有特定格式(例如retailername20230701),可标识对应月份的数据。需实现增量逻辑:若源数据中历史文件有更新,则完全删除该文件的旧数据,插入对应最新load_timestamp的记录;新文件则直接追加。
注意:数据集无唯一键,无法预判哪些字段会变更。
示例数据
sample_table初始数据:
| name | file_name | load_timestamp |
|---|---|---|
| Mark | abc20231201 | 2023-12-01 09:46:48.021448 UTC |
| Antony | xyz20231203 | 2023-12-03 10:26:12.021338 UTC |
| Messi | abc20231201 | 2023-12-01 09:46:48.021448 UTC |
当文件abc20231201的数据更新后,需完全删除该文件的旧记录,替换为最新load_timestamp的记录(如Mark替换为Jack),更新后的源数据:
| Name | file_name | load_timestamp |
|---|---|---|
| Jack | abc20231201 | 2024-01-02 10:16:28.041229 UTC |
| Messi | abc20231201 | 2024-01-02 10:16:28.041229 UTC |
预期结果:
| Name | file_name | load_timestamp |
|---|---|---|
| Antony | xyz20231203 | 2023-12-03 10:26:12.021338 UTC |
| Jack | abc20231201 | 2024-01-02 10:16:28.041229 UTC |
| Messi | abc20231201 | 2024-01-02 10:16:28.041229 UTC |
当前问题
参考DBT官方文档实现了如下代码,但未得到预期结果,无法实现按file_name删除旧记录并插入最新数据:
{{ config( materialized='incremental', partition_by = {'field': 'load_timestamp','data_type': 'timestamp'}, incremental_strategy= 'insert_overwrite' ) }} select * from sample_table {% if is_incremental() %} -- this filter will only be applied on an incremental run -- (uses > to include records whose timestamp occurred since the last run of this model) where load_timestamp >= (select max(load_timestamp) from {{ this }}) {% endif %}
问题分析
原代码存在两个核心问题:
- 增量过滤条件仅加载
load_timestamp大于目标表最大时间戳的数据,无法处理旧文件更新时的旧数据删除,会导致同file_name的新旧数据共存 - 使用
insert_overwrite并按load_timestamp分区覆盖,会直接覆盖整个时间分区的所有数据,而非仅替换对应file_name的旧记录,不符合需求
解决方案
改用delete+insert的增量策略,先删除目标表中本次更新涉及的所有file_name旧记录,再插入最新数据:
{{ config( materialized='incremental', incremental_strategy= 'delete+insert' ) }} with new_data as ( -- 获取本次待加载的所有新数据(含更新的旧文件数据) select * from sample_table ), target_file_names as ( -- 提取本次数据涉及的所有file_name select distinct file_name from new_data ) {% if is_incremental() %} -- 增量运行时,删除目标表中对应file_name的所有旧记录 delete from {{ this }} where file_name in (select file_name from target_file_names); {% endif %} -- 插入本次所有新数据 insert into {{ this }} select * from new_data;
代码说明
- 将
incremental_strategy设置为delete+insert,支持先删除后插入的增量逻辑 - 通过CTE获取本次所有待加载数据,并提取其中的file_name列表
- 增量运行时,删除目标表中该file_name对应的全部旧记录
- 最后插入本次的所有新数据,实现旧文件完全替换、新文件直接追加的需求
内容的提问来源于stack exchange,提问作者Vishal K
相关产品推荐
相关产品推荐

