SQLAlchemy+asyncpg+Procrastinate异步执行时操作冲突错误求助
异步SQLAlchemy结合Procrastinate的并发连接错误解决
问题场景
在使用Procrastinate异步执行SQLAlchemy操作时,运行以下任务代码出现连接冲突错误:
@procrastinate.task(name='application__remove_components') async def remove_applicatoin_components(application_id: int): """ Uninstalls application components. If all application components uninstalled successfully runs post-terminate hooks. """ async with session_maker() as session: application_manager = get_application_manager(session) application = await application_manager.get_application(application_id) manifest: TemplateSchema = load_template(application.manifest) try: await asyncio.gather(*[ application_manager.uninstall_component(application, component) for component in manifest.components if component.enabled ]) except ApplicationComponentUninstallException: await application_manager.set_state_status(application, ApplicationStatuses.error) raise await execute_post_terminate_hooks.defer_async(application_id=application_id)
错误信息
task-executor_1 | /home/app/hub/crud/base.py:34: SAWarning: Usage of the 'Session.add()' operation is not currently supported within the execution stage of the flush process. Results may not be consistent. Consider using alternative event listeners or connection-level operations instead. task-executor_1 | self.session.add(instance) task-executor_1 | ERROR:procrastinate.worker.worker:Job application__remove_components[23373](application_id=78) ended with status: Error, lasted 2.818 s task-executor_1 | Traceback (most recent call last): task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/dialects/postgresql/asyncpg.py", line 739, in commit task-executor_1 | self.await_(self._transaction.commit()) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/util/_concurrency_py3k.py", line 68, in await_only task-executor_1 | return current.driver.switch(awaitable) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/util/_concurrency_py3k.py", line 121, in greenlet_spawn task-executor_1 | value = await result task-executor_1 | File "/usr/local/lib/python3.10/site-packages/asyncpg/transaction.py", line 211, in commit task-executor_1 | await self.__commit() task-executor_1 | File "/usr/local/lib/python3.10/site-packages/asyncpg/transaction.py", line 179, in __commit task-executor_1 | await self._connection.execute(query) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/asyncpg/connection.py", line 317, in execute task-executor_1 | return await self._protocol.query(query, timeout) task-executor_1 | File "asyncpg/protocol/protocol.pyx", line 323, in query task-executor_1 | File "asyncpg/protocol/protocol.pyx", line 707, in asyncpg.protocol.protocol.BaseProtocol._check_state task-executor_1 | asyncpg.exceptions._base.InterfaceError: cannot perform operation: another operation is in progress task-executor_1 | task-executor_1 | The above exception was the direct cause of the following exception: task-executor_1 | task-executor_1 | Traceback (most recent call last): task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/base.py", line 1089, in _commit_impl task-executor_1 | self.engine.dialect.do_commit(self.connection) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/default.py", line 686, in do_commit task-executor_1 | dbapi_connection.commit() task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/dialects/postgresql/asyncpg.py", line 741, in commit task-executor_1 | self._handle_exception(error) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/dialects/postgresql/asyncpg.py", line 682, in _handle_exception task-executor_1 | raise translated_error from error task-executor_1 | sqlalchemy.dialects.postgresql.asyncpg.AsyncAdapt_asyncpg_dbapi.InterfaceError: <class 'asyncpg.exceptions._base.InterfaceError'>: cannot perform operation: another operation is in progress task-executor_1 | task-executor_1 | The above exception was the direct cause of the following exception: task-executor_1 | task-executor_1 | Traceback (most recent call last): task-executor_1 | File "/usr/local/lib/python3.10/site-packages/procrastinate/worker.py", line 231, in run_job task-executor_1 | task_result = await task_result task-executor_1 | File "/home/app/hub/services/procrastinate/tasks/application/terminate_flow.py", line 63, in remove_applicatoin_components task-executor_1 | await asyncio.gather(*[ task-executor_1 | File "/home/app/hub/managers/applications.py", line 380, in uninstall_component task-executor_1 | await self.helm_manager.uninstall_release( task-executor_1 | File "/home/app/hub/managers/helm/manager.py", line 482, in uninstall_release task-executor_1 | await self.event_manager.create(EventSchema( task-executor_1 | File "/home/app/hub/managers/events.py", line 26, in create task-executor_1 | await self.db.create(event.dict()) task-executor_1 | File "/home/app/hub/crud/base.py", line 35, in create task-executor_1 | await self.session.commit() task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/ext/asyncio/session.py", line 582, in commit task-executor_1 | return await greenlet_spawn(self.sync_session.commit) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/util/_concurrency_py3k.py", line 126, in greenlet_spawn task-executor_1 | result = context.throw(*sys.exc_info()) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/orm/session.py", line 1451, in commit task-executor_1 | self._transaction.commit(_to_root=self.future) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/orm/session.py", line 846, in commit task-executor_1 | return self._parent.commit(_to_root=True) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/orm/session.py", line 836, in commit task-executor_1 | trans.commit() task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/base.py", line 2459, in commit task-executor_1 | self._do_commit() task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/base.py", line 2649, in _do_commit task-executor_1 | self._connection_commit_impl() task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/base.py", line 2620, in _connection_commit_impl task-executor_1 | self.connection._commit_impl() task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/base.py", line 1091, in _commit_impl task-executor_1 | self._handle_dbapi_exception(e, None, None, None, None) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/base.py", line 2124, in _handle_dbapi_exception task-executor_1 | util.raise_( task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/util/compat.py", line 210, in raise_ task-executor_1 | raise exception task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/base.py", line 1089, in _commit_impl task-executor_1 | self.engine.dialect.do_commit(self.connection) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/engine/default.py", line 686, in do_commit task-executor_1 | dbapi_connection.commit() task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/dialects/postgresql/asyncpg.py", line 741, in commit task-executor_1 | self._handle_exception(error) task-executor_1 | File "/usr/local/lib/python3.10/site-packages/sqlalchemy/dialects/postgresql/asyncpg.py", line 682, in _handle_exception task-executor_1 | raise translated_error from error task-executor_1 | sqlalchemy.exc.InterfaceError: (sqlalchemy.dialects.postgresql.asyncpg.InterfaceError) <class 'asyncpg.exceptions._base.InterfaceError'>: cannot perform operation: another operation is in progress task-executor_1 | (Background on this error at: https://sqlalche.me/e/14/rvf5)
已尝试的无效方法
- 向engine传入NullPool禁用连接池
- 将连接池大小扩展为pool_size=50,max_overflow=100
- 创建engine、session maker、session的新实例
问题根源
- AsyncSession非协程安全:代码中用
asyncio.gather并发执行多个uninstall_component,所有协程共享同一个外部session。SQLAlchemy的异步session不支持并发操作,多个协程同时在同一个连接上执行事务,会触发连接状态冲突。 - Flush阶段的Session.add操作:SAWarning提示在flush过程中调用
session.add,这种操作会干扰事务的正常执行流程,导致状态不一致。
解决方案
1. 为每个并发任务分配独立Session
修改任务代码,让每个卸载协程使用独立的session,避免共享资源:
@procrastinate.task(name='application__remove_components') async def remove_applicatoin_components(application_id: int): """ Uninstalls application components. If all application components uninstalled successfully runs post-terminate hooks. """ # 单独获取应用信息,使用独立session async with session_maker() as session: application_manager = get_application_manager(session) application = await application_manager.get_application(application_id) manifest: TemplateSchema = load_template(application.manifest) # 定义单个组件卸载的异步函数,内部创建独立session async def uninstall_single_component(component): async with session_maker() as session: app_manager = get_application_manager(session) await app_manager.uninstall_component(application, component) try: await asyncio.gather(*[ uninstall_single_component(component) for component in manifest.components if component.enabled ]) except ApplicationComponentUninstallException: # 错误状态更新也使用独立session async with session_maker() as session: app_manager = get_application_manager(session) await app_manager.set_state_status(application, ApplicationStatuses.error) raise await execute_post_terminate_hooks.defer_async(application_id=application_id)
2. 修正flush阶段的Session.add操作
检查crud/base.py中的create方法,调整session.add的执行时机,避免在flush过程中调用:
# 修改crud/base.py的create方法 async def create(self, obj_in): db_obj = self.model(**obj_in) self.session.add(db_obj) # 先执行flush再提交,确保add操作不在其他事务流程中 await self.session.flush() await self.session.commit() await self.session.refresh(db_obj) return db_obj
如果是在事件监听函数中调用session.add,改用connection级别的操作(如直接执行INSERT语句),或者调整监听事件的触发时机,避开flush阶段。
3. 确保Procrastinate Worker的资源隔离
检查Procrastinate的worker配置,确保每个worker协程有独立的数据库上下文,避免跨任务共享连接或session实例。
内容的提问来源于stack exchange,提问作者Vadim Zabolotniy
相关产品推荐
相关产品推荐

