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

Azure ADO CICD中pytest无法解析Airflow变量的解决求助

问题:Azure ADO CICD中pytest无法解析Airflow变量

背景

在Azure ADO CICD管道运行pytest test_create_fivetran_job.py时出错,该测试脚本用于验证create_fivetran_job.py。

功能说明

create_fivetran_job.py从Snowflake配置表读取配置,通过API调用在Fivetran创建刷新任务。该脚本由Airflow MWAA DAG调度,所需的AUTH_URL、Fivetran登录凭证、Snowflake凭证均存储在Airflow变量中。

当前问题

Python脚本和DAG运行正常,但pytest执行失败,无法解析Airflow变量。已确认Airflow变量与yml文件内容正确,脚本可被DAG正常触发。

错误信息

______________ ERROR collecting tests/test_create_fivetran_job.py ______________
tests/test_create_fivetran_job.py:43: in <module>
    from ..dags.scripts import create_fivetran_job
dags/scripts/create_fivetran_job.py:19: in <module>
    AUTH_URL = Variable.get("fivetran_auth_url")
/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/airflow/models/variable.py:140: in get
    raise KeyError(f'Variable {key} does not exist')
E   KeyError: 'Variable fivetran_auth_url does not exist' 

捕获的标准输出

[2024-02-05 20:01:09,116] {variable.py:274} ERROR - Unable to retrieve variable from secrets backend (MetastoreBackend). Checking subsequent secrets backend.
Traceback (most recent call last):
File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1911, in _execute_context
cursor, statement, parameters, context
File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/engine/default.py", line 736, in do_execute
cursor.execute(statement, parameters)
sqlite3.OperationalError: no such table: variable

    The above exception was the direct cause of the following exception:

    Traceback (most recent call last):
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/airflow/models/variable.py", line 267, in get_variable_from_secrets
        var_val = secrets_backend.get_variable(key=key)
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/airflow/utils/session.py", line 70, in wrapper
        return func(*args, session=session, **kwargs)
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/airflow/secrets/metastore.py", line 64, in get_variable
        var_value = session.query(Variable).filter(Variable.key == key).first()
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/orm/query.py", line 2824, in first
        return self.limit(1)._iter().first()
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/orm/query.py", line 2919, in _iter
        execution_options={"_sa_orm_load_options": self.load_options},
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/orm/session.py", line 1717, in execute
        result = conn._execute_20(statement, params or {}, execution_options)
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1710, in _execute_20
        return meth(self, args_10style, kwargs_10style, execution_options)
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/sql/elements.py", line 335, in _execute_on_connection
        self, multiparams, params, execution_options
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1587, in _execute_clauseelement
        cache_hit=cache_hit,
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1954, in _execute_context
        e, statement, parameters, cursor, context
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 2135, in _handle_dbapi_exception
        sqlalchemy_exception, with_traceback=exc_info[2], from_=e
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/util/compat.py", line 211, in raise_
        raise exception
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/engine/base.py", line 1911, in _execute_context
        cursor, statement, parameters, context
      File "/opt/hostedtoolcache/Python/3.7.17/x64/lib/python3.7/site-packages/sqlalchemy/engine/default.py", line 736, in do_execute
        cursor.execute(statement, parameters)
    sqlalchemy.exc.OperationalError: (sqlite3.OperationalError) no such table: variable
    [SQL: SELECT variable.val AS variable_val, variable.id AS variable_id, variable."key" AS variable_key, variable.description AS variable_description, variable.is_encrypted AS variable_is_encrypted 
    FROM variable 
    WHERE variable."key" = ?
     LIMIT ? OFFSET ?]
    [parameters: ('fivetran_auth_url', 1, 0)]

解决方法

方法1:初始化Airflow测试环境

利用Airflow测试依赖创建本地元数据表,模拟生产环境的变量存储:

  1. 安装测试依赖:
pip install apache-airflow[testing]
  1. 在测试文件中添加全局初始化fixture:
import pytest
from airflow.utils.db import initdb
from airflow.models import Variable
from airflow.utils.session import create_session

@pytest.fixture(scope="session", autouse=True)
def setup_airflow_env():
    # 创建Airflow元数据表
    initdb()
    # 预加载测试用变量
    with create_session() as session:
        Variable.set("fivetran_auth_url", "test_auth_url", session=session)
        # 按需添加其他变量

方法2:通过环境变量传递变量

Airflow支持从环境变量读取变量,格式为AIRFLOW_VAR_<变量名>(变量名大写,下划线保留):

  1. 在ADO CICD管道中配置环境变量,例如将fivetran_auth_url对应设置为AIRFLOW_VAR_FIVETRAN_AUTH_URL
  2. 或在测试文件中手动设置:
import os
os.environ["AIRFLOW_VAR_FIVETRAN_AUTH_URL"] = "test_auth_url"

方法3:重构脚本实现依赖注入

修改create_fivetran_job.py,允许变量通过参数传入,避免模块级直接调用Variable.get:

# create_fivetran_job.py
from airflow.models import Variable

def main(auth_url=None, fivetran_creds=None, snowflake_creds=None):
    AUTH_URL = auth_url or Variable.get("fivetran_auth_url")
    # 剩余业务逻辑...

if __name__ == "__main__":
    main()

测试时直接传入测试值:

# test_create_fivetran_job.py
from dags.scripts import create_fivetran_job

def test_create_job():
    create_fivetran_job.main(
        auth_url="test_url",
        fivetran_creds={"user": "test_user", "api_key": "test_key"},
        snowflake_creds={"account": "test_account"}
    )
    # 执行断言逻辑

方法4:Mock替换Variable.get

使用pytest的mock功能直接覆盖Variable.get方法,返回预设测试值:

# test_create_fivetran_job.py
from unittest.mock import patch

def test_script_import():
    with patch('airflow.models.Variable.get') as mock_get:
        mock_get.side_effect = lambda key: {
            "fivetran_auth_url": "test_url",
            # 映射其他需要的变量
        }[key]
        # 导入脚本或执行测试逻辑
        from dags.scripts import create_fivetran_job
        # 后续测试断言

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 17:19:50