Airflow PostgresOperator如何获取API调用传入的conf参数
Airflow PostgresOperator 获取API传入的conf参数值
场景说明
通过API触发DAG时传入包含seller_id的配置参数,API调用命令如下:
curl --location 'http://localhost:8080/api/v1/dags/postgres_operator_dag/dagRuns' \ --header 'Content-Type: application/json' \ --data '{ "conf": { "seller_id": "test_seller" } }'
失败尝试
用户尝试以下三种方式均未成功:
- 在
parameters中传入模板字符串{{ dag_run.conf["seller_id"] }},任务返回无数据,但硬编码值可正常运行; - 在SQL语句中直接使用带引号的模板字符串,系统将其视为普通字符串,无法正确查询;
- 在SQL语句中直接使用模板字符串,出现psycopg2语法错误。
正确实现方式
方法一:使用PostgresOperator的parameters参数(推荐,避免SQL注入)
Airflow的PostgresOperator支持通过parameters传递参数,结合Jinja2模板获取dag_run.conf中的值,同时让psycopg2自动处理参数转义,彻底规避SQL注入风险。
示例代码:
from airflow.providers.postgres.operators.postgres import PostgresOperator postgres_task = PostgresOperator( task_id="query_seller_data", postgres_conn_id="your_postgres_conn", sql=""" SELECT * FROM your_table WHERE seller_id = %(seller_id)s; """, parameters={ "seller_id": "{{ dag_run.conf['seller_id'] }}" }, dag=dag )
核心逻辑:SQL语句使用%(seller_id)s作为占位符,parameters字典中的值通过Jinja模板渲染获取dag_run.conf里的seller_id,Airflow会自动完成模板渲染,再由psycopg2处理参数绑定。
方法二:直接在SQL中使用Jinja模板并处理字符串格式
如果要直接在SQL里嵌入模板变量,需要确保Jinja正确渲染值并处理字符串引号。可以使用Jinja的string过滤器保证渲染结果为字符串类型:
示例代码:
postgres_task = PostgresOperator( task_id="query_seller_data", postgres_conn_id="your_postgres_conn", sql=""" SELECT * FROM your_table WHERE seller_id = '{{ dag_run.conf['seller_id'] | string }}'; """, dag=dag )
注意:这种方式如果seller_id包含特殊字符(比如单引号),可能引发SQL注入或语法错误,因此方法一更安全可靠。
关键注意点
- 确保DAG的
render_template_as_native_obj设置为True(Airflow 2.x+支持),让Jinja渲染后的参数保持原类型,避免不必要的字符串转义问题:from airflow import DAG with DAG( dag_id="postgres_operator_dag", render_template_as_native_obj=True, # 其他DAG参数(schedule_interval、start_date等) ) as dag: # 任务定义... - 如果
dag_run.conf中的seller_id可能不存在,建议添加默认值,避免模板渲染失败:# 在parameters中添加默认值 parameters={ "seller_id": "{{ dag_run.conf.get('seller_id', 'default_seller') }}" } # 或者在SQL模板中添加默认值 WHERE seller_id = '{{ dag_run.conf.get('seller_id', 'default_seller') | string }}';
内容的提问来源于stack exchange,提问作者Arjunsingh
相关产品推荐
相关产品推荐

