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

Airflow模板变量在PostgresHook中解析失败的问题求助

Airflow模板变量在自定义PostgresHook代码中无法替换的问题

问题现象

在Airflow中使用{{ds}}、{{execution_date.strftime('%Y-%m-%d')}}等模板变量时,PostgresOperator中替换正常,但在自定义的PostgresHook代码里直接读取YAML配置文件时,模板变量完全不生效,执行retention类型任务时触发以下错误:

psycopg2.errors.UndefinedColumn: column "y" does not exist
LINE 1: ...((date_trunc('month','{{execution_date.strftime('%Y-%m-%d')}...

自定义PostgresHook代码

def prc_mymys_update(procedure: str, type_agg: str):
    with PostgresHook(postgres_conn_id=CONNECTION_ID_GP).get_conn() as conn:
        with conn.cursor() as cur:
            with open(URL_YML_2,"r", encoding="utf-8") as f:
                ya_2 = yaml.safe_load(f)
                yml_mymts_2 = ya_2['type_agg']
                query_pg = ""
                if yml_mymts_2[0]['type_agg_name'] == "day" and type_agg == "day":
                    sql_1 = yml_mymts_2[0]['sql']
                    query_pg = f"""{sql_1}"""
                elif yml_mymts_2[1]['type_agg_name'] == "retention" and type_agg == "retention":
                    sql_2 = yml_mymts_2[1]['sql']
                    query_pg = f"""{sql_2}"""
                elif yml_mymts_2[2]['type_agg_name'] == "mau" and type_agg == "mau":
                    sql_3 = yml_mymts_2[2]['sql']
                    query_pg = f"""{sql_3}"""
                cur.execute(query_pg)
                dates_list = cur.fetchall()
                for date_res in dates_list:
                    cur.execute(
                        "select from {}(%(date)s::date);".format(procedure),
                        {"date": date_res[0].strftime("%Y-%m-%d")},
                    )
    conn.close()

YAML配置文件

type_agg:
  - type_agg_name: day
    sql: select calendar_date from entertainment_dds.v_calendar where calendar_date between '{{ds}}'::date - interval '7 days' and '{{ds}}'::date - 1 order by 1 desc
  - type_agg_name: retention
    sql: SELECT t.date::date AS date FROM generate_series((date_trunc('month','{{execution_date.strftime('%Y-%m-%d')}}'::date) - interval '11 month'), date_trunc('month','{{execution_date.strftime('%Y-%m-%d')}}'::date) , '1 month'::interval) t(date) order by 1 asc
  - type_agg_name: mau
    sql: select dt::date date_ from generate_series('{{execution_date.strftime('%Y-%m-%d')}}'::date - interval '7 days', '{{execution_date.strftime('%Y-%m-%d')}}'::date - interval '1 days', interval '1 days') dt order by 1 asc

问题根源

  1. 模板渲染机制限制:Airflow的模板变量替换是在Operator层面由Jinja2引擎自动处理的,直接通过open()读取YAML文件时,内容不会经过Airflow的模板渲染流程,因此{{ds}}等变量会被当作原始字符串传入SQL。
  2. YAML单引号嵌套冲突:YAML配置中strftime('%Y-%m-%d')的单引号与外层包裹SQL的单引号冲突,导致YAML解析时字符串被截断,生成的SQL语法错误,触发"column y does not exist"报错。

解决方法

步骤1:修复YAML中的单引号问题

将YAML中strftime里的单引号替换为双引号,避免嵌套冲突:

type_agg:
  - type_agg_name: day
    sql: select calendar_date from entertainment_dds.v_calendar where calendar_date between '{{ds}}'::date - interval '7 days' and '{{ds}}'::date - 1 order by 1 desc
  - type_agg_name: retention
    sql: SELECT t.date::date AS date FROM generate_series((date_trunc('month','{{execution_date.strftime("%Y-%m-%d")}}'::date) - interval '11 month'), date_trunc('month','{{execution_date.strftime("%Y-%m-%d")}}'::date) , '1 month'::interval) t(date) order by 1 asc
  - type_agg_name: mau
    sql: select dt::date date_ from generate_series('{{execution_date.strftime("%Y-%m-%d")}}'::date - interval '7 days', '{{execution_date.strftime("%Y-%m-%d")}}'::date - interval '1 days', interval '1 days') dt order by 1 asc

步骤2:在自定义代码中实现模板变量替换

有两种可行方案:

方案A:手动传入Airflow上下文变量并替换

修改自定义函数,接收execution_date和ds参数(通过PythonOperator的provide_context=True传递),然后手动替换SQL中的模板变量:

def prc_mymys_update(procedure: str, type_agg: str, execution_date=None, ds=None):
    with PostgresHook(postgres_conn_id=CONNECTION_ID_GP).get_conn() as conn:
        with conn.cursor() as cur:
            with open(URL_YML_2,"r", encoding="utf-8") as f:
                ya_2 = yaml.safe_load(f)
                yml_mymts_2 = ya_2['type_agg']
                query_pg = ""
                target_date = execution_date.strftime("%Y-%m-%d")
                
                if yml_mymts_2[0]['type_agg_name'] == "day" and type_agg == "day":
                    query_pg = yml_mymts_2[0]['sql'].replace("{{ds}}", ds)
                elif yml_mymts_2[1]['type_agg_name'] == "retention" and type_agg == "retention":
                    query_pg = yml_mymts_2[1]['sql'].replace("{{execution_date.strftime(\"%Y-%m-%d\")}}", target_date)
                elif yml_mymts_2[2]['type_agg_name'] == "mau" and type_agg == "mau":
                    query_pg = yml_mymts_2[2]['sql'].replace("{{execution_date.strftime(\"%Y-%m-%d\")}}", target_date)
                
                cur.execute(query_pg)
                dates_list = cur.fetchall()
                for date_res in dates_list:
                    cur.execute(
                        "select from {}(%(date)s::date);".format(procedure),
                        {"date": date_res[0].strftime("%Y-%m-%d")},
                    )

然后在PythonOperator中启用上下文传递:

PythonOperator(
    task_id='update_retention_data',
    python_callable=prc_mymys_update,
    op_kwargs={'procedure': 'your_procedure_name', 'type_agg': 'retention'},
    provide_context=True,
    dag=dag
)

方案B:使用Jinja2引擎渲染YAML文件

将YAML文件放在Airflow的templates目录下,通过Jinja2引擎加载并渲染,自动替换模板变量:

from jinja2 import Environment, FileSystemLoader

def prc_mymys_update(procedure: str, type_agg: str, execution_date=None, ds=None):
    # 初始化Jinja2环境,指定模板目录(需替换为你的templates路径)
    env = Environment(loader=FileSystemLoader('/opt/airflow/dags/templates'))
    # 加载YAML模板文件
    template = env.get_template('your_config.yml')
    # 渲染YAML,传入Airflow上下文变量
    rendered_yml_content = template.render(ds=ds, execution_date=execution_date)
    # 解析渲染后的YAML
    ya_2 = yaml.safe_load(rendered_yml_content)
    
    with PostgresHook(postgres_conn_id=CONNECTION_ID_GP).get_conn() as conn:
        with conn.cursor() as cur:
            yml_mymts_2 = ya_2['type_agg']
            query_pg = ""
            if yml_mymts_2[0]['type_agg_name'] == "day" and type_agg == "day":
                query_pg = yml_mymts_2[0]['sql']
            elif yml_mymts_2[1]['type_agg_name'] == "retention" and type_agg == "retention":
                query_pg = yml_mymts_2[1]['sql']
            elif yml_mymts_2[2]['type_agg_name'] == "mau" and type_agg == "mau":
                query_pg = yml_mymts_2[2]['sql']
            
            cur.execute(query_pg)
            dates_list = cur.fetchall()
            for date_res in dates_list:
                cur.execute(
                    "select from {}(%(date)s::date);".format(procedure),
                    {"date": date_res[0].strftime("%Y-%m-%d")},
                )

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 12:50:22