如何构建无重复的增量模型?解决增量聚合重复条目问题
解决dbt增量模型重复聚合条目的问题
问题原因
当前代码的核心问题是增量运行时仅聚合新增的源数据并直接追加到目标表,如果同一业务日期(DATE(DATE_TIME))有多批次数据流入(比如同一天不同小时的新数据),每次增量运行都会生成该日期的一条聚合记录,最终目标表中同一日期会出现多条重复条目,而非单一的完整聚合结果。
解决方案:使用Merge策略重写增量模型
通过dbt的merge增量策略,我们可以对目标表中已存在的日期记录进行覆盖,确保每个日期仅保留最新的完整聚合结果。同时,为了处理可能的旧日期补数场景,我们会重新拉取新增数据涉及的所有业务日期的完整源数据进行聚合,保证结果准确性。
修改后的代码如下:
{{ config( materialized='incremental', incremental_strategy='merge', -- 指定使用merge策略 unique_key='DATE' -- 唯一键设为业务日期,确保每个日期仅一条记录 )}} WITH WEBSITE_MAIN_PAGE_VIEWS AS ( SELECT DATE_TIME, POST_VISID_HIGH, POST_VISID_LOW, VISIT_NUM, VISIT_START_TIME_GMT, POST_PAGENAME, POST_EVENT_LIST, _DMLTIME FROM {{ ref('stg_adobe_hit_data_poc') }} WHERE 1 = 1 {% if is_incremental() %} -- 取出新增数据涉及的所有业务日期的完整源数据,保证聚合准确性 DATE(DATE_TIME) IN ( SELECT DISTINCT DATE(DATE_TIME) FROM {{ ref('stg_adobe_hit_data_poc') }} WHERE _DMLTIME > ( SELECT COALESCE(MAX(_DMLTIME), TO_TIMESTAMP_TZ('1900-01-01 00:00:00')) FROM {{ this }} ) ) {% endif %} ), new_aggregations AS ( SELECT 'PO' AS BRAND, DATE(DATE_TIME) AS DATE, '*' AS BASE_MARKET, COUNT(DISTINCT CONCAT(POST_VISID_HIGH, POST_VISID_LOW, VISIT_NUM, VISIT_START_TIME_GMT)) AS TOTAL_WEBSITE_VISITS, COUNT( DISTINCT CASE WHEN POST_PAGENAME = 'po:en_gb' THEN CONCAT(POST_VISID_HIGH, POST_VISID_LOW, VISIT_NUM, VISIT_START_TIME_GMT) END ) AS HOME_PAGE_VISITS, MAX(_DMLTIME) AS _DMLTIME FROM WEBSITE_MAIN_PAGE_VIEWS GROUP BY DATE(DATE_TIME) ) SELECT * FROM new_aggregations
关键修改说明
- 配置增量策略:在
config中指定incremental_strategy='merge'和unique_key='DATE',dbt会自动根据DATE字段匹配目标表记录,用新的聚合结果覆盖旧记录。 - 调整增量数据范围:不再仅取
_DMLTIME大于目标表最大值的源数据,而是取出这些新增数据涉及的所有业务日期的完整源数据,重新计算聚合,避免因补数导致的结果不准确。 - 统一聚合逻辑:全量和增量模式共用同一聚合逻辑,保证结果一致性。
简化场景(无旧日期补数)
如果可以确定源表不会出现旧业务日期的补数(新增数据的DATE_TIME均为最新业务日期),可以简化增量过滤条件,直接聚合新增源数据后merge:
{{ config( materialized='incremental', incremental_strategy='merge', unique_key='DATE' )}} WITH WEBSITE_MAIN_PAGE_VIEWS AS ( SELECT DATE_TIME, POST_VISID_HIGH, POST_VISID_LOW, VISIT_NUM, VISIT_START_TIME_GMT, POST_PAGENAME, POST_EVENT_LIST, _DMLTIME FROM {{ ref('stg_adobe_hit_data_poc') }} WHERE 1 = 1 {% if is_incremental() %} AND _DMLTIME > ( SELECT COALESCE(MAX(_DMLTIME), TO_TIMESTAMP_TZ('1900-01-01 00:00:00')) FROM {{ this }} ) {% endif %} ), new_aggregations AS ( SELECT 'PO' AS BRAND, DATE(DATE_TIME) AS DATE, '*' AS BASE_MARKET, COUNT(DISTINCT CONCAT(POST_VISID_HIGH, POST_VISID_LOW, VISIT_NUM, VISIT_START_TIME_GMT)) AS TOTAL_WEBSITE_VISITS, COUNT( DISTINCT CASE WHEN POST_PAGENAME = 'po:en_gb' THEN CONCAT(POST_VISID_HIGH, POST_VISID_LOW, VISIT_NUM, VISIT_START_TIME_GMT) END ) AS HOME_PAGE_VISITS, MAX(_DMLTIME) AS _DMLTIME FROM WEBSITE_MAIN_PAGE_VIEWS GROUP BY DATE(DATE_TIME) ) SELECT * FROM new_aggregations
这种情况下,merge策略会自动将新增的聚合记录与目标表中同日期的记录合并,覆盖旧值,避免重复。
内容的提问来源于stack exchange,提问作者Abiodun Adeoye
相关产品推荐
相关产品推荐

