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

如何在Airflow中运行DBT模型并在运行时传递appId参数

实现方案:通过Airflow传递appId参数给DBT模型

1. 让DBT模型支持接收appId参数

在你的DBT模型文件(如abc.sql)中,使用var()函数获取外部传入的appId参数,过滤对应数据:

-- 单个appId的情况
select * from target_table
where app_id = '{{ var("appId") }}'

-- 多个appId的情况(对应你期望的appIdList逻辑)
select * from target_table
where app_id in ({{ "'" + "','".join(var("appIdList")) + "'" }})

2. 修改Airflow的KubernetesPodOperator任务

通过DBT CLI的--vars参数,将Airflow变量中的appId值传递给DBT任务,有两种适配场景:

场景一:传递单个appId

abc_task = KubernetesPodOperator(
    namespace='etl',
    image=f'blotout/dbt-analytics:{TAG_DBT_ANALYTICS}',
    cmds=["/usr/local/bin/dbt"],
    arguments=[
        'run',
        '--models', 'abc_task',
        '--vars', '{"appId": "{{ var.value.appId }}"}'
    ],
    env_vars=env_var,
    name="abc_task",
    configmaps=['awskey'],
    task_id="abc_task",
    get_logs=True,
    dag=dag,
    is_delete_operator_pod=True,
)

场景二:传递appId列表

abc_task = KubernetesPodOperator(
    namespace='etl',
    image=f'blotout/dbt-analytics:{TAG_DBT_ANALYTICS}',
    cmds=["/usr/local/bin/dbt"],
    arguments=[
        'run',
        '--models', 'abc_task',
        '--vars', '{"appIdList": {{ var.value.appIdList | tojson }}}'
    ],
    env_vars=env_var,
    name="abc_task",
    configmaps=['awskey'],
    task_id="abc_task",
    get_logs=True,
    dag=dag,
    is_delete_operator_pod=True,
)

注:使用tojson过滤器可将Airflow的列表变量直接转为DBT能解析的JSON格式

3. 传入appId参数的两种方式

  • 提前配置Airflow变量:在Airflow UI的「Admin > Variables」中添加appId(或appIdList)变量并设置对应值
  • 手动触发DAG时传参:触发DAG时在「Conf」栏输入JSON格式参数,比如:
    {"appId": "your-target-app-id"}
    
    或列表形式:
    {"appIdList": ["app-001", "app-002"]}
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 09:28:11