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

构建不更新已有记录的dbt增量模型:获取用户最小时间戳

解决dbt增量模型中保留user_id最小updated_at且避免不必要更新的问题

现有代码的核心问题

  1. Unique Key设置错误:原模型将version_id设为唯一键,但你的需求是每个user_id仅保留一条最小updated_at的记录,因此唯一键应改为user_id,否则会导致同一user_id的多条记录被重复插入。
  2. 增量过滤逻辑不严谨:仅通过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的记录,初始化目标表。
  • 增量运行:
    1. 先获取目标表中已有的user_id和对应的最小时间戳。
    2. 仅筛选两种记录:
      • 目标表中从未出现过的user_id,提取其最小updated_at插入。
      • 已存在的user_id,但源数据中出现了比目标表记录更早的激活时间(比如旧记录刚被标记为激活),此时会更新目标表中的对应记录。
    3. 当源数据中新增某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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 05:10:21