使用unittest.mock模拟Airflow Livy Operator连接时遇连接未定义错误
解决Airflow Livy Operator测试中连接未定义的问题
问题描述
使用unittest.mock的@mock.patch.dict装饰器设置环境变量AIRFLOW_CONN_LIVY_HOOK以模拟Airflow Livy Operator连接时,运行测试报错:
airflow.exceptions.AirflowNotFoundException: The conn_id 'AIRFLOW_CONN_LIVY_HOOK' isn't defined
问题原因
- 连接ID命名规则错误:Airflow中环境变量命名格式为
AIRFLOW_CONN_<CONN_ID>,但在Operator中指定conn_id时,只需使用<CONN_ID>部分,而非完整的环境变量名。 - mock参数格式错误:
@mock.patch.dict的环境变量参数需以字典形式传入,而非关键字参数。 - Airflow连接缓存未清空:Airflow会缓存已加载的连接,即使mock了环境变量,也可能读取缓存中的旧值。
解决方案
1. 修正mock装饰器的参数格式
将环境变量以字典形式传入@mock.patch.dict,同时设置clear=False(避免清空所有环境变量,影响其他测试依赖)。
2. 修正LivyOperator的conn_id取值
将livy_conn_id的值改为"LIVY_HOOK",对应环境变量AIRFLOW_CONN_LIVY_HOOK。
3. 清空Airflow连接缓存
在测试的setUp方法中调用connections.clear(),确保Airflow读取最新的mock环境变量。
修正后的完整代码示例
import unittest from unittest import mock from airflow import DAG from airflow.providers.apache.livy.operators.livy import LivyOperator from airflow.utils import timezone import requests_mock # 假设spark_args已定义 spark_args = ["--conf", "spark.executor.instances=2"] class TestLivyOperator(unittest.TestCase): def setUp(self): super().setUp() # 清空Airflow连接缓存,确保读取mock的环境变量 from airflow.models import connections connections.clear() self.dag = DAG( dag_id="test_livy", default_args={ "owner": "xyz", "start_date": timezone.datetime(2022, 8, 16), }, ) @mock.patch.dict( "os.environ", {"AIRFLOW_CONN_LIVY_HOOK": "http://www.google.com"}, clear=False ) @requests_mock.mock() def test_payload(self, mock_request): # 模拟Livy API的响应 mock_request.post("http://www.google.com/batches", json={"id": 1}, status_code=201) task = LivyOperator( task_id="task_1", class_name="com.precious.myClass", executor_memory="512m", executor_cores=3, arg=spark_args, livy_conn_id="LIVY_HOOK", # 使用正确的conn_id dag=self.dag, ) # 执行任务测试 task.execute(context={}) # 验证请求是否符合预期 assert mock_request.called assert mock_request.call_count == 1
额外说明
- 若测试类中多个方法需要使用该mock,可将
@mock.patch.dict装饰在类上,但需确保clear=False,避免影响其他测试的环境变量。 - 若使用Airflow 2.x版本,
connections.clear()的导入路径可能为from airflow.models.connection import Connection; Connection.clear(),请根据实际版本调整。
内容的提问来源于stack exchange,提问作者Priyanshu Sharma
相关产品推荐
相关产品推荐

