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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 07:28:12