如何用dbt增量模型小时级维护商品可用性状态
问题描述
我有一个销售多种商品的网站,存在一张记录所有商品点击行为的clicks表,表结构及数据如下:
| click_id | product_id |
|---|---|
| 1 | 1234 |
| 2 | 0023 |
| 3 | 1234 |
场景说明
- 商品并非每日都可售;
- 仅能通过商品出现在
clicks表(即有用户点击)判断其存在。
此前通过SELECT DISTINCT product_id FROM clicks可查询曾可用的商品,但该全量查询数据量过大,无法每日执行。
需求目标
需要维护一份商品历史列表,实现:
- 每小时查询当日被点击的商品;
- 新增从未出现过的商品并标记
is_available=1; - 根据当日点击情况更新商品的可用性状态。
尝试的dbt增量模型代码(未成功)
{{ config( materialized='incremental', unique_key=['product_id'], ) }} WITH online_products AS ( SELECT DISTINCT product_id, 1 AS is_available, CURRENT_DATE() AS updated_at FROM clicks WHERE partition_date >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR) ) SELECT DISTINCT product_id function_to_update(is_online), function_to_update(updated_at) FROM {{ this }} current {% if is_incremental() %} FULL OUTER JOIN online_products USING product_id {% endif %}
解决方案
你的核心问题在于增量模型的逻辑未处理好新旧数据的合并规则,以及is_available状态的更新逻辑。以下是修正后的dbt增量模型实现:
修正后的代码
{{ config( materialized='incremental', unique_key='product_id', -- 单个唯一键可简化写法 incremental_strategy='merge' -- 明确指定merge策略,适配主流数据仓库 ) }} WITH hourly_clicked_products AS ( -- 提取最近1小时内被点击的唯一商品,标记为可用状态 SELECT DISTINCT product_id, 1 AS is_available, CURRENT_TIMESTAMP() AS updated_at -- 用时间戳记录更精确的更新时间 FROM clicks WHERE partition_date = CURRENT_DATE() AND click_timestamp >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR) ) SELECT COALESCE(current.product_id, hourly.product_id) AS product_id, -- 状态更新逻辑:当前小时有点击则标记为1,历史商品无点击则保留原有状态 -- 若需求为"当日无点击则标记为不可用",可改为 COALESCE(hourly.is_available, 0) CASE WHEN hourly.product_id IS NOT NULL THEN 1 ELSE current.is_available END AS is_available, COALESCE(hourly.updated_at, current.updated_at) AS updated_at FROM {{ this }} current {% if is_incremental() %} -- 增量运行时,关联当前小时的点击商品数据 FULL OUTER JOIN hourly_clicked_products hourly ON current.product_id = hourly.product_id {% else %} -- 全量初始化时,加载所有历史点击过的商品(首次运行用全量逻辑) RIGHT JOIN ( SELECT DISTINCT product_id, 1 AS is_available, CURRENT_TIMESTAMP() AS updated_at FROM clicks ) hourly ON current.product_id = hourly.product_id {% endif %}
关键逻辑说明
增量策略与关联方式
- 使用
merge增量策略,通过product_id作为唯一键匹配新旧数据 - 增量运行时用
FULL OUTER JOIN覆盖三类场景:新商品(仅在小时点击表中)、历史商品(仅在现有模型表中)、状态更新的商品(两边都存在) - 全量初始化时直接查询所有历史点击商品,确保首次运行能加载完整的商品列表
- 使用
可用性状态控制
- 优先以当前小时的点击记录为准,有点击则标记为可用(
1) - 无点击的历史商品保留原有状态,若需强制当日无点击则标记为不可用,可修改
CASE语句的ELSE分支为0
- 优先以当前小时的点击记录为准,有点击则标记为可用(
时间范围优化
- 同时过滤日期和小时级时间戳,缩小数据扫描范围,提升查询效率
- 用
CURRENT_TIMESTAMP()替代日期记录更新时间,更精准反映状态变化时点
额外优化建议
- 如果
clicks表按小时分区,可直接使用partition_date >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR)进一步减少数据扫描量 - 若需追踪商品每日可用性历史,可新增
record_date字段,每日生成一条状态记录而非覆盖原有数据
内容的提问来源于stack exchange,提问作者Antoine DEMATTEO
相关产品推荐
相关产品推荐

