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
问题根源
- 模板渲染机制限制:Airflow的模板变量替换是在Operator层面由Jinja2引擎自动处理的,直接通过
open()读取YAML文件时,内容不会经过Airflow的模板渲染流程,因此{{ds}}等变量会被当作原始字符串传入SQL。 - 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
相关产品推荐
相关产品推荐

