AWS Lambda Python3.8异步SQLAlchemy触发IllegalStateChangeError问题
解决SQLAlchemy AsyncSession在Lambda中的IllegalStateChangeError错误
环境信息
- AWS Lambda(Python3.8)
- SQLAlchemy 2.0.10
- asyncpg 0.28.0
问题代码
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession engine_async = create_async_engine( "postgresql+asyncpg://postgres:password@host_name.rds.amazonaws.com/postgres") def get_data(event, context): # this is the lambda handler return asyncio.get_event_loop().run_until_complete(get_data_async(event, context)) async def get_datas_async(event, context): org_id = event['org_id'] async_session = sessionmaker(bind=engine_async, future=True, class_=AsyncSession) async with async_session() as session1, async_session() as session2, async_session() as session3: users_count_query = select(func.count(User.user_id)).filter_by(org_id=org_id) org_name_query = select(Organization.name).filter_by(org_id=org_id) other_data_query = select( func.count(Transaction.transaction_id).label('transactions_count'), func.sum(Transaction.total_cost).label('total_cost'), func.sum(Transaction.total_time).label('total_time')).where(Transaction.org_id == org_id) tasks = [ session1.execute(users_count_query), session2.execute(other_data_query), session3.execute(org_name_query) ] results = await asyncio.gather(*tasks) userss_count, data, org_name = results data = data.fetchone() uesrs_count = users_count.scalar_one() org_name = org_name.scalar_one() return { "statusCode": 200, "users_count": users_count, "transactions_count": data.transactions_count, "total_cost": 0 if not data.total_cost else data.total_cost, "total_time": data.total_time, "org_name": org_name, }
错误信息
{ "errorMessage": "Method 'close()' can't be called here; method '_connection_for_bind()' is already in progress and this would cause an unexpected state change to <SessionTransactionState.CLOSED: 5>", "errorType": "IllegalStateChangeError", "stackTrace": [ " File \"/var/task/app.py\", line 131, in get_data\n return asyncio.get_event_loop().run_until_complete(get_data_async(event, context))\n", " File \"/var/lang/lib/python3.8/asyncio/base_events.py\", line 616, in run_until_complete\n return future.result()\n", " File \"/var/task/app.py\", line 175, in get_data_async\n return {\n", " File \"/var/task/sqlalchemy/ext/asyncio/session.py\", line 859, in __aexit__\n await asyncio.shield(task)\n", " File \"/var/task/sqlalchemy/ext/asyncio/session.py\", line 840, in close\n await greenlet_spawn(self.sync_session.close)\n", " File \"/var/task/sqlalchemy/util/_concurrency_py3k.py\", line 154, in greenlet_spawn\n result = context.switch(*args, **kwargs)\n", " File \"/var/task/sqlalchemy/orm/session.py\", line 2382, in close\n self._close_impl(invalidate=False)\n", " File \"/var/task/sqlalchemy/orm/session.py\", line 2424, in _close_impl\n transaction.close(invalidate)\n", " File \"<string>\", line 2, in close\n", " File \"/var/task/sqlalchemy/orm/state_changes.py\", line 121, in _go\n raise sa_exc.IllegalStateChangeError(\n" ] }
解决方案
问题根源
错误源于同时创建多个AsyncSession并在async with块中并行执行操作,当块退出时多个会话同时尝试关闭,导致SQLAlchemy内部会话事务状态冲突。SQLAlchemy的异步会话设计支持在单个会话中并行执行多个查询,无需为每个查询单独创建会话。
修正后的代码
import asyncio from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import sessionmaker from sqlalchemy import select, func # 请确保User、Organization、Transaction是已定义的模型类 from your_models_module import User, Organization, Transaction # 全局初始化引擎和会话工厂,避免每次请求重复创建 engine_async = create_async_engine( "postgresql+asyncpg://postgres:password@host_name.rds.amazonaws.com/postgres" ) AsyncSessionLocal = sessionmaker( bind=engine_async, future=True, class_=AsyncSession ) def get_data(event, context): return asyncio.get_event_loop().run_until_complete(get_data_async(event, context)) async def get_data_async(event, context): org_id = event['org_id'] # 使用单个AsyncSession处理所有并行查询 async with AsyncSessionLocal() as session: # 定义查询语句 users_count_query = select(func.count(User.user_id)).filter_by(org_id=org_id) org_name_query = select(Organization.name).filter_by(org_id=org_id) other_data_query = select( func.count(Transaction.transaction_id).label('transactions_count'), func.sum(Transaction.total_cost).label('total_cost'), func.sum(Transaction.total_time).label('total_time') ).where(Transaction.org_id == org_id) # 在单个会话中并行执行所有查询任务 tasks = [ session.execute(users_count_query), session.execute(other_data_query), session.execute(org_name_query) ] users_count_result, other_data_result, org_name_result = await asyncio.gather(*tasks) # 解析查询结果 users_count = users_count_result.scalar_one() other_data = other_data_result.fetchone() org_name = org_name_result.scalar_one() return { "statusCode": 200, "users_count": users_count, "transactions_count": other_data.transactions_count, "total_cost": 0 if not other_data.total_cost else other_data.total_cost, "total_time": other_data.total_time, "org_name": org_name, }
关键修正点
- 将
sessionmaker定义为全局变量,避免每次Lambda调用重复创建会话工厂 - 用单个AsyncSession替代多个会话,SQLAlchemy异步会话可安全处理多个并发查询
- 修正了原代码中的变量名拼写错误(
get_datas_async、userss_count、uesrs_count) - 简化会话上下文管理,避免多会话关闭时的状态冲突
内容的提问来源于stack exchange,提问作者Rony Tesler
相关产品推荐
相关产品推荐

