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

如何使用SnowflakeOperator生成含dag_id等信息的动态Query Tag

如何用SnowflakeOperator动态设置包含dag_id、task_id、run_id的Query Tag

要实现动态Query Tag,核心是利用Airflow的Jinja模板变量——SnowflakeOperator的session_parameters参数支持模板渲染,可直接引用Airflow内置的任务上下文变量。

修改后的代码示例

with DAG(
    "test_querytags",
    description="Refreshes entities shared by all projects",
    default_args=default_args,
    schedule='15 8 * * *',
    max_active_runs=1,
    catchup=False,
    default_view="graph",
    tags=["admin","common","calendar","daily"],
) as dag:

    t_Load_Data = SnowflakeOperator(
        task_id="SF__Load_Data",
        sql="CALL SANDBOX_OR.RAW.LOAD_CITY_2_P('SANDBOX_OR')",
        # 用Jinja模板动态拼接Query Tag
        session_parameters={
            "QUERY_TAG": "{{ dag.dag_id }} | {{ task.task_id }} | {{ run_id }}"
        }
    )

关键说明

  • {{ dag.dag_id }}:获取当前DAG的唯一标识
  • {{ task.task_id }}:获取当前任务的ID
  • {{ run_id }}:获取当前DAG运行实例的ID(包含时间戳和运行类型,比如手动触发会带manual__前缀)

生成的Query Tag会自动填充对应的值,在Snowflake查询历史中能直接看到这条查询所属的DAG、任务和运行实例,大幅提升排查效率。

可选优化:处理特殊字符

如果run_id包含Snowflake Query Tag不允许的特殊字符(比如冒号:),可用Jinja过滤器替换:

session_parameters={
    "QUERY_TAG": "{{ dag.dag_id }} | {{ task.task_id }} | {{ run_id | replace(':', '-') }}"
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 20:33:18