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测试依赖创建本地元数据表,模拟生产环境的变量存储:
- 安装测试依赖:
pip install apache-airflow[testing]
- 在测试文件中添加全局初始化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_<变量名>(变量名大写,下划线保留):
- 在ADO CICD管道中配置环境变量,例如将
fivetran_auth_url对应设置为AIRFLOW_VAR_FIVETRAN_AUTH_URL - 或在测试文件中手动设置:
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
相关产品推荐
相关产品推荐

