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_FOOisn'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
核心问题
HttpSensor的
http_conn_id不是环境变量,是Airflow Connection的ID
Airflow的HttpSensor会根据http_conn_id去Airflow内置的Connections配置中查找对应连接,而非读取环境变量。你把环境变量名当成conn_id传入,自然找不到匹配的连接。任务进程隔离导致环境变量不共享
Airflow每个任务都在独立进程执行,init_vars任务里设置的环境变量只能在自身进程生效,无法传递给后续的task_is_api_active任务。直接调用
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

