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

如何使用异步SQLAlchemy Engine实现Alembic数据库迁移?

异步SQLAlchemy + Alembic迁移失败问题排查与解决

问题背景

将SQLAlchemy PostgreSQL驱动转换为异步版本后,使用异步引擎执行Alembic数据库迁移时出现异常:迁移流程显示完成,但数据初始化(Seeding)环节报错,同时出现协程未等待的RuntimeWarning。

错误日志

2023-12-24 14:36:22 Starting entrypoint.sh
2023-12-24 14:36:29 Database is up and running
2023-12-24 14:36:29 Generating migrations
2023-12-24 14:36:38 Generating /app/alembic/versions/9a4735888d4b_initial_migration.py ...  done
2023-12-24 14:36:41 Running migrations
2023-12-24 14:36:45 Migration completed successfully.
2023-12-24 14:36:45 Seeding with test user
2023-12-24 14:36:49 An error occurred while seeding the expressions: AsyncConnection context has not been started and object has not been awaited.
2023-12-24 14:36:50 Inside start_server function
2023-12-24 14:36:50 Starting ngrok
Authtoken saved to configuration file: /root/.config/ngrok/ngrok.yml                                
2023-12-24 14:36:22 wait-for-it.sh: waiting 60 seconds for db:5432
2023-12-24 14:36:29 wait-for-it.sh: db:5432 is available after 7 seconds
2023-12-24 14:36:38 INFO  [alembic.runtime.migration] Context impl PostgresqlImpl.
2023-12-24 14:36:38 INFO  [alembic.runtime.migration] Will assume transactional DDL.
2023-12-24 14:36:45 INFO  [alembic.runtime.migration] Context impl PostgresqlImpl.
2023-12-24 14:36:45 INFO  [alembic.runtime.migration] Will assume transactional DDL.
2023-12-24 14:36:45 INFO  [alembic.runtime.migration] Running upgrade  -> 9a4735888d4b, Initial migration
2023-12-24 14:36:49 /usr/local/lib/python3.11/site-packages/sqlalchemy/orm/session.py:775: RuntimeWarning: coroutine 'AsyncConnection.close' was never awaited
2023-12-24 14:36:49   conn.close()

相关代码文件

alembic/env.py

from logging.config import fileConfig
import asyncio
from sqlalchemy.ext.asyncio import create_async_engine
from sqlalchemy.pool import NullPool
from alembic import context
from database.database_config import Base, db_url
from services.utils import logger
import traceback

config = context.config
fileConfig(config.config_file_name)
target_metadata = Base.metadata

if db_url:
    config.set_main_option("sqlalchemy.url", db_url)

def do_run_migrations(connection):
    try:
        context.configure(
            connection=connection,
            target_metadata=target_metadata
        )

        with context.begin_transaction():
            context.run_migrations()
    except Exception as e:
        logger.error(traceback.format_exc())
        raise

async def run_async_migrations():
    connectable = create_async_engine(db_url, poolclass=NullPool)

    async with connectable.connect() as connection:
        await connection.run_sync(do_run_migrations)

    await connectable.dispose()

def run_migrations_online():
    asyncio.run(run_async_migrations())

run_migrations_online()

database_config.py

from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import sessionmaker
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession

from services.utils import logger
import traceback
from config.env_var import *

DB_USER = os.getenv('DB_USER')
DB_PASSWORD = os.getenv('DB_PASSWORD')
DB_HOST = os.getenv('DB_HOST')
DB_NAME = os.getenv('DB_NAME')

Base = declarative_base()


db_url = f'postgresql+asyncpg://{DB_USER}:{DB_PASSWORD}@{DB_HOST}:5432/{DB_NAME}'

try:
    engine = create_async_engine(db_url, echo=True)
except Exception as e:
    logger.info(f"Error creating database engine: {e}")
    logger.info(traceback.format_exc())
    raise

AsyncSessionLocal = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)

async def get_db():
    db = AsyncSessionLocal()
    try:
        yield db
    except Exception as e:
        logger.info(f"Failed with db_url: {db_url}")
        logger.info(f"Database session error: {e}")
        logger.info(traceback.format_exc())
        raise
    finally:
        await db.close()

init_db.py

import asyncio
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from sqlalchemy.orm import sessionmaker

from database.models import *
from database.enums import *
from database.database_config import Base, engine, db_url

async def create_tables():
    # Use the async engine from your database configuration
    async_engine = create_async_engine(db_url)

    # Asynchronous table creation
    async with async_engine.begin() as conn:
        await conn.run_sync(Base.metadata.create_all)

if __name__ == "__main__":
    asyncio.run(create_tables())

解决方案

1. 修复Alembic迁移的事务逻辑

修改alembic/env.py中的do_run_migrations函数,避免异步上下文与Alembic事务管理冲突:

def do_run_migrations(connection):
    try:
        context.configure(
            connection=connection,
            target_metadata=target_metadata,
            # 让Alembic为每个迁移单独处理事务
            transaction_per_migration=True,
            compare_type=True
        )
        # 移除手动事务上下文,由run_sync自动处理
        context.run_migrations()
    except Exception as e:
        logger.error(traceback.format_exc())
        raise

原因:with context.begin_transaction()在异步run_sync环境中会导致事务上下文混乱,改用transaction_per_migration=True更适配异步迁移场景。

2. 修正数据初始化的异步调用

确保Seeding逻辑完全在异步上下文中执行,使用AsyncSession而非直接操作连接:

# 示例Seeding代码(替换为实际业务逻辑)
async def seed_test_user():
    async with AsyncSessionLocal() as session:
        async with session.begin():
            test_user = User(username="test", email="test@example.com")
            session.add(test_user)
        await session.commit()

# 在启动脚本中通过asyncio.run执行
asyncio.run(seed_test_user())

原因:报错提示异步连接未启动,说明Seeding代码可能在同步上下文调用了异步数据库操作,必须全程使用异步Session并等待所有异步方法。

3. 消除协程未等待警告

检查所有数据库操作代码,确保:

  • 所有异步方法(如session.close()、session.commit())都添加await关键字
  • 禁止在同步函数中直接调用异步Session/Connection的方法

4. 优化init_db.py的引擎复用

init_db.py无需重复创建引擎,直接复用database_config.py中已定义的引擎:

async def create_tables():
    # 复用统一配置的异步引擎
    async with engine.begin() as conn:
        await conn.run_sync(Base.metadata.create_all)

原因:重复创建引擎会造成资源浪费,复用全局引擎更符合规范。

验证步骤

  1. 重新生成迁移:alembic revision --autogenerate -m "fix async migration"
  2. 执行迁移:alembic upgrade head
  3. 运行数据初始化脚本,确认无报错
  4. 检查日志,确认无coroutine was never awaited警告

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 15:21:05