Airflow单元测试添加SubDag时出现slot_pool表不存在错误求助
问题:添加SubDag后单元测试出现SQLite表不存在错误
我用tox+pytest对Mock DAG做单元测试,原本正常,加了SubDag后出现SQL错误:
sqlalchemy.exc.OperationalError: (sqlite3.OperationalError) no such table: slot_pool
补充:我有个自定义库测试时也会创建SubDag,但没触发这个错误。
我的Mock代码:
@pytest.fixture def mock_dag(): dag = DAG( "my_dag", default_args={ "start_date": datetime(2021, 7, 15), }, ) sub_dag = DAG( dag_id="my_dag.my_subdag", params=dag.params, start_date=datetime(2021, 7, 15), schedule_interval=dag.schedule_interval, default_args=dag.default_args, template_searchpath=dag.template_searchpath, user_defined_macros=dag.user_defined_macros, ) sub_dag.is_subdag = True DummyOperator(task_id="test", dag=sub_dag) SubDagOperator( task_id="my_subdag", subdag=sub_dag, default_args=dag.default_args, dag=dag, ) return dag
完整堆栈信息:
return __inject_session_to_func(func, session_type, *args, **kwargs) .tox/py37/lib/python3.7/site-packages/airflow/utils/db.py:72: in __inject_session_to_func return func(*args, **kwargs) .tox/py37/lib/python3.7/site-packages/airflow/utils/decorators.py:98: in wrapper result = func(*args, **kwargs) .tox/py37/lib/python3.7/site-packages/airflow/operators/subdag_operator.py:84: in __init__ .filter(Pool.pool == self.pool) .tox/py37/lib/python3.7/site-packages/sqlalchemy/orm/query.py:3429: in first ret = list(self[0:1]) .tox/py37/lib/python3.7/site-packages/sqlalchemy/orm/query.py:3203: in __getitem__ return list(res) .tox/py37/lib/python3.7/site-packages/sqlalchemy/orm/query.py:3535: in __iter__ return self._execute_and_instances(context) .tox/py37/lib/python3.7/site-packages/sqlalchemy/orm/query.py:3560: in _execute_and_instances result = conn.execute(querycontext.statement, self._params) .tox/py37/lib/python3.7/site-packages/sqlalchemy/engine/base.py:1011: in execute return meth(self, multiparams, params) .tox/py37/lib/python3.7/site-packages/sqlalchemy/sql/elements.py:298: in _execute_on_connection return connection._execute_clauseelement(self, multiparams, params) .tox/py37/lib/python3.7/site-packages/sqlalchemy/engine/base.py:1130: in _execute_clauseelement distilled_params, .tox/py37/lib/python3.7/site-packages/sqlalchemy/engine/base.py:1317: in _execute_context e, statement, parameters, cursor, context .tox/py37/lib/python3.7/site-packages/sqlalchemy/engine/base.py:1511: in _handle_dbapi_exception sqlalchemy_exception, with_traceback=exc_info[2], from_=e .tox/py37/lib/python3.7/site-packages/sqlalchemy/util/compat.py:182: in raise_ raise exception .tox/py37/lib/python3.7/site-packages/sqlalchemy/engine/base.py:1277: in _execute_context cursor, statement, parameters, context _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ self = <sqlalchemy.dialects.sqlite.pysqlite.SQLiteDialect_pysqlite object at 0x106c77e80>, cursor = <sqlite3.Cursor object at 0x11efdc110> statement = 'SELECT slot_pool.id AS slot_pool_id, slot_pool.pool AS slot_pool_pool, slot_pool.slots AS slot_pool_slots, slot_pool....iption AS slot_pool_description \nFROM slot_pool \nWHERE slot_pool.slots = ? AND slot_pool.pool = ?\n LIMIT ? OFFSET ?' parameters = (1, 'default_pool', 1, 0), context = <sqlalchemy.dialects.sqlite.base.SQLiteExecutionContext object at 0x11efd69b0> def do_execute(self, cursor, statement, parameters, context=None): > cursor.execute(statement, parameters) E sqlalchemy.exc.OperationalError: (sqlite3.OperationalError) no such table: slot_pool E [SQL: SELECT slot_pool.id AS slot_pool_id, slot_pool.pool AS slot_pool_pool, slot_pool.slots AS slot_pool_slots, slot_pool.description AS slot_pool_description E FROM slot_pool E WHERE slot_pool.slots = ? AND slot_pool.pool = ? E LIMIT ? OFFSET ?] E [parameters: (1, 'default_pool', 1, 0)] E (Background on this error at: http://sqlalche.me/e/13/e3q8) .tox/py37/lib/python3.7/site-packages/sqlalchemy/engine.default.p
解决方案
从堆栈信息可以看出,SubDagOperator初始化时会查询slot_pool表验证池配置,但你的测试环境没有初始化Airflow元数据库表。以下是几种解决方式:
1. 初始化测试用元数据库
使用Airflow自带的命令在测试前创建所有元数据表:
airflow db init
或者在pytest中添加一个session级别的fixture自动初始化:
import pytest from airflow.utils.db import initdb @pytest.fixture(scope="session") def init_airflow_db(): initdb()
然后让mock_dag fixture依赖这个初始化fixture:
@pytest.fixture def mock_dag(init_airflow_db): # 原代码保持不变
2. Mock池验证逻辑绕过数据库依赖
如果不想初始化全量数据库,可以用pytest mock掉SubDagOperator中查询池的逻辑:
from unittest.mock import patch @pytest.fixture def mock_dag(): dag = DAG( "my_dag", default_args={ "start_date": datetime(2021, 7, 15), }, ) sub_dag = DAG( dag_id="my_dag.my_subdag", params=dag.params, start_date=datetime(2021, 7, 15), schedule_interval=dag.schedule_interval, default_args=dag.default_args, template_searchpath=dag.template_searchpath, user_defined_macros=dag.user_defined_macros, ) sub_dag.is_subdag = True DummyOperator(task_id="test", dag=sub_dag) # Mock Pool查询,避免直接访问数据库 with patch('airflow.operators.subdag_operator.Pool') as mock_pool: mock_pool.query.filter.return_value.first.return_value = None subdag_op = SubDagOperator( task_id="my_subdag", subdag=sub_dag, default_args=dag.default_args, dag=dag, ) return dag
3. 复用自定义库的测试配置
你的自定义库没触发错误,大概率是它的测试环境已经做了数据库初始化或Mock了池验证逻辑。可以对比自定义库的conftest.py等配置文件,找到对应的初始化或Mock代码直接复用。
内容的提问来源于stack exchange,提问作者pavbagel
相关产品推荐
相关产品推荐

