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

如何通过Airflow执行DBT增量模型的指定日期范围刷新?

实现按指定日期范围刷新dbt增量模型的方案

问题背景

我有一个采用insert_override策略的dbt增量模型,每小时通过Airflow DAG增量运行,删除并更新过去2天的数据。当源schema变更或派生表存在数据缺失时,会触发临时DAG并携带--full-refresh标志执行全量刷新,但仅单天数据缺失时全量刷新成本过高。当前全量刷新需修改Airflow变量,配置如下:

{
  "_comment": "NOTE: after changing values here, refresh the DAG from the Airflow UI.",
  "dbt_parser_list": [
    "--select", "model_name", "--full-refresh"
  ],
  "dbt_commands": [
    "run"
  ]
}

希望通过临时DAG传入日期,仅刷新该日期之后的数据(例如传入2024-04-01即可刷新4月及以后的数据,传入默认起始日期2020-01-01则执行全量刷新),了解到--vars参数可用于代码通用化,但不清楚具体用法,求实现方法。


解决方案

1. 改造dbt增量模型逻辑

修改dbt模型代码,通过vars参数接收传入的起始日期,调整数据过滤逻辑:

{{ config(
    materialized='incremental',
    incremental_strategy='insert_override',
    unique_key='your_unique_identifier' -- 必须指定唯一键,用于insert_override匹配要覆盖的行
) }}

WITH source_data AS (
    SELECT * 
    FROM {{ source('your_source_schema', 'your_source_table') }}
    -- 根据传入的起始日期过滤数据
    WHERE 
        {% if is_incremental() %}
            -- 增量运行场景:优先使用传入的start_date,否则默认取过去2天
            {% if var('start_date', none) %}
                your_date_column >= '{{ var('start_date') }}'
            {% else %}
                your_date_column >= DATEADD(day, -2, CURRENT_DATE())
            {% endif %}
        {% else %}
            -- 非增量运行场景(首次运行或指定全量起始日期):取start_date之后的数据
            your_date_column >= '{{ var('start_date', '2020-01-01') }}'
        {% endif %}
)

SELECT * FROM source_data

核心说明:

  • var('start_date', '2020-01-01'):获取传入的变量,第二个参数为默认值
  • is_incremental():判断当前是增量还是全量运行模式
  • unique_key:insert_override策略依赖该参数识别需要覆盖的行

2. 调整Airflow临时DAG的变量配置

删除原配置中的--full-refresh,改用--vars传递日期参数:

{
  "_comment": "NOTE: after changing values here, refresh the DAG from the Airflow UI.",
  "dbt_parser_list": [
    "--select", "model_name",
    "--vars", "{\"start_date\": \"2024-04-01\"}"
  ],
  "dbt_commands": [
    "run"
  ]
}

若需全量刷新,只需将start_date设为默认值2020-01-01即可。

3. 让临时DAG支持动态传入日期参数

为避免每次修改Airflow变量,可给临时DAG添加交互式日期输入参数(以Airflow PythonOperator为例):

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
import subprocess

def run_dbt_with_date(start_date='2020-01-01'):
    dbt_cmd = [
        "dbt", "run",
        "--select", "model_name",
        "--vars", f'{{"start_date": "{start_date}"}}'
    ]
    subprocess.run(dbt_cmd, check=True)

with DAG(
    'dbt_partial_refresh_dag',
    start_date=datetime(2024, 1, 1),
    catchup=False,
    params={
        'start_date': {
            'type': 'string',
            'description': '起始日期,格式YYYY-MM-DD,默认2020-01-01为全量刷新',
            'default': '2020-01-01'
        }
    }
) as dag:
    run_dbt_task = PythonOperator(
        task_id='run_dbt_partial_refresh',
        python_callable=run_dbt_with_date,
        op_kwargs={'start_date': '{{ params.start_date }}'}
    )

触发DAG时,可直接在Airflow UI中输入目标起始日期,无需修改变量。

4. 测试验证

在命令行先验证逻辑是否正确:

# 刷新2024-04-01及以后的数据
dbt run --select model_name --vars '{"start_date": "2024-04-01"}'

# 执行全量刷新
dbt run --select model_name --vars '{"start_date": "2020-01-01"}'

内容的提问来源于stack exchange,提问作者nitika sharma

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 12:32:36