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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 21:07:17