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

Airflow中HttpSensor无法识别环境变量定义的conn_id问题排查

问题描述

我写了一个尝试连接HTTP端点的Airflow DAG,通过PythonOperator定义环境变量AIRFLOW_VAR_FOO,但HttpSensor识别不了这个conn_id,报错“The conn_id AIRFLOW_VAR_FOO isn't defined”。我试着直接调用init_vars()也没解决,以下是DAG代码和完整错误信息,请问问题出在哪?

报错信息

The conn_id AIRFLOW_VAR_FOO isn't defined

原始DAG代码

import os
import json
import pprint
import datetime
import requests

from airflow.models import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.sftp.operators.sftp import SFTPOperator
from airflow.providers.sftp.sensors.sftp import SFTPSensor
from airflow.utils.dates import days_ago
from airflow.models import Variable
from airflow.sensors.http_sensor import HttpSensor
from airflow.hooks.base_hook import BaseHook

def init_vars():
    os.environ['AIRFLOW_VAR_FOO'] = "https://mywebxxx.net/"
    print(os.environ['AIRFLOW_VAR_FOO'])

with DAG(
         dag_id='request_test',
         schedule_interval=None,
         start_date=days_ago(2)) as dag:

    init_vars = PythonOperator(task_id="init_vars",
                                  python_callable=init_vars)

    task_is_api_active = HttpSensor(
        task_id='is_api_active',
        http_conn_id='AIRFLOW_VAR_FOO',
        endpoint='post'
    )

    get_data = PythonOperator(task_id="get_data",
                                  python_callable=get_data)

    init_vars >> task_is_api_active

完整错误日志

File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/connection.py", line 379, in get_connection_from_secrets
    raise AirflowNotFoundException(f"The conn_id `{conn_id}` isn't defined")
airflow.exceptions.AirflowNotFoundException: The conn_id `AIRFLOW_VAR_FOO` isn't defined
[2022-11-04 10:32:41,720] {taskinstance.py:1551} INFO - Marking task as FAILED. dag_id=request_test, task_id=is_api_active, execution_date=20221104T103235, start_date=20221104T103240, end_date=20221104T103241
[2022-11-04 10:32:42,628] {local_task_job.py:149} INFO - Task exited with return code 1

编辑后的DAG代码

import os
import json
import pprint
import datetime
import requests

from airflow.models import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.sftp.operators.sftp import SFTPOperator
from airflow.providers.sftp.sensors.sftp import SFTPSensor
from airflow.utils.dates import days_ago
from airflow.models import Variable
from airflow.sensors.http_sensor import HttpSensor
from airflow.hooks.base_hook import BaseHook

def init_vars():
    os.environ['AIRFLOW_VAR_FOO'] = "https://mywebxxx.net/"
    print(os.environ['AIRFLOW_VAR_FOO'])

with DAG(
         dag_id='request_test',
         schedule_interval=None,
         start_date=days_ago(2)) as dag:

    init_vars = PythonOperator(task_id="init_vars",
                                  python_callable=init_vars)

    call init_vars()
    
    task_is_api_active = HttpSensor(
        task_id='is_api_active',
        http_conn_id='AIRFLOW_VAR_FOO',
        endpoint='post'
    )

    get_data = PythonOperator(task_id="get_data",
                                  python_callable=get_data)

    task_is_api_active

问题原因及解决方案

核心问题

  1. HttpSensor的http_conn_id不是环境变量,是Airflow Connection的ID
    Airflow的HttpSensor会根据http_conn_id去Airflow内置的Connections配置中查找对应连接,而非读取环境变量。你把环境变量名当成conn_id传入,自然找不到匹配的连接。

  2. 任务进程隔离导致环境变量不共享
    Airflow每个任务都在独立进程执行,init_vars任务里设置的环境变量只能在自身进程生效,无法传递给后续的task_is_api_active任务。

  3. 直接调用init_vars()的时机错误
    编辑后的代码里直接调用init_vars(),这是在DAG解析阶段执行的,设置的环境变量只存在于解析进程中,任务执行时的独立进程根本访问不到。

正确解决方案

方案1:创建Airflow HTTP Connection(推荐)

在Airflow UI的Admin > Connections中创建HTTP类型连接:

  • Conn ID:自定义名称,比如my_http_conn
  • Host:填写https://mywebxxx.net/
  • 其他字段(如端口、认证信息)按需补充

修改HttpSensor参数:

task_is_api_active = HttpSensor(
    task_id='is_api_active',
    http_conn_id='my_http_conn',  # 使用创建的conn_id
    endpoint='post'
)

方案2:使用Airflow Variable(无需创建Connection场景)

先在Airflow UI的Admin > Variables中添加变量:

  • Key:FOO
  • Value:https://mywebxxx.net/

自定义传感器读取变量并检查API:

from airflow.sensors.base import BaseSensorOperator
import requests

class CustomHttpSensor(BaseSensorOperator):
    def __init__(self, endpoint, **kwargs):
        super().__init__(**kwargs)
        self.endpoint = endpoint

    def poke(self, context):
        base_url = Variable.get("FOO")
        full_url = f"{base_url}/{self.endpoint}"
        try:
            response = requests.get(full_url)
            return response.status_code == 200
        except Exception as e:
            self.log.error(f"API访问失败: {str(e)}")
            return False

# 使用自定义传感器
task_is_api_active = CustomHttpSensor(
    task_id='is_api_active',
    endpoint='post'
)

方案3:PythonOperator直接检查API(无需传感器重试逻辑场景)

如果不需要传感器的自动重试机制,直接用PythonOperator调用requests:

def check_api_status():
    base_url = Variable.get("FOO")
    full_url = f"{base_url}/post"
    response = requests.get(full_url)
    if response.status_code != 200:
        raise Exception(f"API不可用,状态码: {response.status_code}")

task_is_api_active = PythonOperator(
    task_id='is_api_active',
    python_callable=check_api_status
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 23:15:52