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

使用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

问题原因

  1. 连接ID命名规则错误:Airflow中环境变量命名格式为AIRFLOW_CONN_<CONN_ID>,但在Operator中指定conn_id时,只需使用<CONN_ID>部分,而非完整的环境变量名。
  2. mock参数格式错误:@mock.patch.dict的环境变量参数需以字典形式传入,而非关键字参数。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 22:55:18