如何使用异步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)
原因:重复创建引擎会造成资源浪费,复用全局引擎更符合规范。
验证步骤
- 重新生成迁移:
alembic revision --autogenerate -m "fix async migration" - 执行迁移:
alembic upgrade head - 运行数据初始化脚本,确认无报错
- 检查日志,确认无
coroutine was never awaited警告
内容的提问来源于stack exchange,提问作者Ari
相关产品推荐
相关产品推荐

