使用PostgresOperator时Airflow模板未渲染XCom值的问题
解决PostgresOperator中XCom值未正确渲染的问题
嘿,我刚好碰到过一模一样的坑!问题出在Airflow对PostgresOperator和BashOperator的模板处理逻辑不一样,我来给你拆解清楚:
核心原因
PostgresOperator只会对**sql参数指向的模板文件内容**做Jinja渲染,而不会处理你传入params里的字符串模板。你把"{{ ti.xcom_pull(key='return_value') }}"作为params.status_msg的值传递时,它只是一个普通字符串,SQL模板里的{{ params.status_msg }}只会原样输出这个字符串,不会解析里面的Jinja语法。
而BashOperator不同,它的bash_command参数本身会被当作Jinja模板处理,所以里面的{{ ... }}能自动渲染成实际值。
最简洁的解决方案:直接在SQL模板里获取XCom
把获取XCom的逻辑移到SQL模板文件里,不需要通过params传递模板字符串,这样最省心:
修改后的SQL模板(update_monitoring_table_status_msg.sql)
UPDATE table SET status_msg = '{{ ti.xcom_pull(task_ids="t1", key="return_value") }}' WHERE script_name = '{{ params.script_name }}'
修改后的DAG代码
from airflow.models import Variable from airflow.operators.python import PythonOperator from airflow.operators.postgres_operator import PostgresOperator from datetime import datetime tmpl_search_path = Variable.get('sql_path') dag = DAG( 'test_dag', description='', schedule_interval='@daily', template_searchpath=tmpl_search_path, start_date=datetime(2017, 3, 20), catchup=False ) # Returns a string t1 = PythonOperator(task_id='t1', python_callable=someCallable, dag=dag) update_monitoring_table_last_run = PostgresOperator( task_id='update_monitoring_table_last_run', sql='update_monitoring_table_last_run.sql', postgres_conn_id='conn_id', params={"script_name": t1.task_id}, dag=dag ) update_monitoring_table_status_msg = PostgresOperator( task_id='update_monitoring_table_status_msg', sql='update_monitoring_table_status_msg.sql', postgres_conn_id='conn_id', params={"script_name": t1.task_id}, # 不再传递status_msg参数 dag=dag ) t1 >> update_monitoring_table_last_run >> update_monitoring_table_status_msg
备选方案:手动渲染params里的模板字符串
如果你一定要通过params传递,也可以手动用Airflow的Jinja工具渲染模板字符串,不过步骤会麻烦一点:
from airflow.models import Variable from airflow.operators.python import PythonOperator from airflow.operators.postgres_operator import PostgresOperator from airflow.templates.jinja_template import JinjaTemplate from datetime import datetime tmpl_search_path = Variable.get('sql_path') dag = DAG( 'test_dag', description='', schedule_interval='@daily', template_searchpath=tmpl_search_path, start_date=datetime(2017, 3, 20), catchup=False ) # Returns a string t1 = PythonOperator(task_id='t1', python_callable=someCallable, dag=dag) update_monitoring_table_last_run = PostgresOperator( task_id='update_monitoring_table_last_run', sql='update_monitoring_table_last_run.sql', postgres_conn_id='conn_id', params={"script_name": t1.task_id}, dag=dag ) # 手动渲染status_msg的模板字符串 status_msg_template = JinjaTemplate("{{ ti.xcom_pull(task_ids='t1', key='return_value') }}") update_monitoring_table_status_msg = PostgresOperator( task_id='update_monitoring_table_status_msg', sql='update_monitoring_table_status_msg.sql', postgres_conn_id='conn_id', params={ "script_name": t1.task_id, "status_msg": status_msg_template.render(context=dag.get_template_context()) }, dag=dag ) t1 >> update_monitoring_table_last_run >> update_monitoring_table_status_msg
不过这个方案需要注意上下文的时效性,不如第一个方案稳定。
内容的提问来源于stack exchange,提问作者Insworn
相关产品推荐
相关产品推荐

