构建不更新已有记录的dbt增量模型:获取用户最小时间戳
解决dbt增量模型中保留user_id最小updated_at且避免不必要更新的问题
现有代码的核心问题
- Unique Key设置错误:原模型将
version_id设为唯一键,但你的需求是每个user_id仅保留一条最小updated_at的记录,因此唯一键应改为user_id,否则会导致同一user_id的多条记录被重复插入。 - 增量过滤逻辑不严谨:仅通过
updated_at > (SELECT MAX(this.updated_at) FROM {{ this }})筛选增量数据,会漏掉源数据中已存在user_id但更早的激活记录(比如某条旧记录的is_activated从false改为true),同时无法精准控制仅处理需要更新的user_id。
修改后的模型代码
{{ config( materialized='incremental', unique_key='user_id' ) }} WITH existing_user_min_times AS ( {% if is_incremental() %} -- 增量运行时,获取目标表中已有的user_id及其最小updated_at SELECT user_id, updated_at AS min_updated_at FROM {{ this }} {% else %} -- 全量运行时,返回空结果集 SELECT NULL::VARCHAR AS user_id, NULL::TIMESTAMP AS min_updated_at {% endif %} ), source_activated_records AS ( SELECT version_id, user_id, updated_at, -- 按user_id分组,标记出最小updated_at的记录 ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY updated_at ASC) AS row_num FROM {{ ref('versions') }} WHERE is_activated = TRUE ), filtered_records AS ( SELECT version_id, user_id, updated_at FROM source_activated_records WHERE row_num = 1 {% if is_incremental() %} AND ( -- 仅处理目标表中不存在的user_id user_id NOT IN (SELECT user_id FROM existing_user_min_times) -- 或目标表中已存在该user_id,但源数据中有更早的激活时间 OR updated_at < (SELECT min_updated_at FROM existing_user_min_times WHERE user_id = source_activated_records.user_id) ) {% endif %} ) SELECT * FROM filtered_records
逻辑说明
- 全量运行:会扫描所有
is_activated = TRUE的源数据,为每个user_id提取最小updated_at的记录,初始化目标表。 - 增量运行:
- 先获取目标表中已有的
user_id和对应的最小时间戳。 - 仅筛选两种记录:
- 目标表中从未出现过的
user_id,提取其最小updated_at插入。 - 已存在的
user_id,但源数据中出现了比目标表记录更早的激活时间(比如旧记录刚被标记为激活),此时会更新目标表中的对应记录。
- 目标表中从未出现过的
- 当源数据中新增某
user_id的晚于现有记录的激活数据时,该user_id的最小时间戳不会变化,因此不会被筛选出来,避免了不必要的更新。
- 先获取目标表中已有的
额外优化建议
如果源表数据量极大,可在source_activated_records的WHERE条件中增加增量时间范围过滤,比如:
{% if is_incremental() %} AND updated_at >= (SELECT COALESCE(MIN(min_updated_at) - INTERVAL '7 days', '1970-01-01'::TIMESTAMP) FROM existing_user_min_times) {% endif %}
这样可以缩小扫描范围,提升增量运行效率(假设不会出现比已捕获最小时间早7天以上的新激活记录,可根据业务调整时间范围)。
内容的提问来源于stack exchange,提问作者kimi
相关产品推荐
相关产品推荐

