如何在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
相关产品推荐
相关产品推荐

