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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 14:36:20