Python asyncio中pubsub回调调用asyncpg报数据库操作冲突如何解决
asyncpg报错
cannot perform operation: another operation is in progress 解决方案 核心原因
asyncpg的连接和连接池与创建它的异步事件循环强绑定,跨事件循环调用、或者同一个连接同时执行多个异步操作,都会触发该报错。
现有代码问题点
- 事件循环不统一:主事件循环
loop1创建了数据库连接池pool,但在Pub/Sub回调中使用Asyncio.run()每次都会启动新的事件循环,跨循环调用连接池直接触发冲突 - 新建无效事件循环:
receive_message函数中创建了新的事件循环loop但没有实际使用,属于冗余逻辑 - 同步回调调用异步逻辑方式错误:Google Pub/Sub的订阅回调默认在独立的同步线程池执行,直接用
asyncio.run()处理协程既会导致事件循环冲突,也存在线程安全问题 - 代码细节错误:
write_to_database函数定义了两个入参,但调用时只传了一个;async with pool.acquire()上下文会自动释放连接,finally中重复释放会导致连接池状态异常;subscriber.subscribe传入的是subscription_id而非拼接好的subscription_path,会导致订阅失败
修复代码
import asyncio import asyncpg import google.cloud.pubsub_v1 # 全局保存主事件循环与连接池,供回调使用 main_loop = None pool = None async def write_to_database(pool): query = """insert into table1 values('a','b');""" async with pool.acquire() as con: await con.execute(query) async def pubsub_callback(message): """Process the incoming message.""" try: print((f'Received ID:{message.message_id} ' f'PUBTIME:{message.publish_time} ' f'ATTEMPT:{message.delivery_attempt} ' f'Data: {message.data}')) await asyncio.sleep(1) await write_to_database(pool) print(f'Message ID:{message.message_id} Processed.') except ValueError: print(f'Message ID:{message.message_id} was not processed') finally: message.ack() async def receive_message(): await asyncio.sleep(.5) subscriber = google.cloud.pubsub_v1.SubscriberClient() project_id = 'project_id' subscription_id = 'subscription_id' with subscriber: subscription_path = subscriber.subscription_path(project_id, subscription_id) print(f'Listening for messages on {subscription_path}') def create_pubsub_callback_task(message): """同步回调中线程安全提交异步任务到主事件循环""" asyncio.run_coroutine_threadsafe(pubsub_callback(message), main_loop) streaming_pull_future = subscriber.subscribe(subscription_path, callback=create_pubsub_callback_task) try: streaming_pull_future.result(timeout=None) except (TimeoutError, KeyboardInterrupt): streaming_pull_future.cancel() if __name__ == "__main__": url = "替换为你的postgres连接串" main_loop = asyncio.get_event_loop() pool = main_loop.run_until_complete(asyncpg.create_pool(url)) main_loop.run_until_complete(receive_message())
修复说明
- 统一使用同一个主事件循环处理所有异步逻辑,包括数据库操作,避免跨循环调用连接池
- 使用
asyncio.run_coroutine_threadsafe在Pub/Sub的同步回调线程中,线程安全地提交异步任务到主事件循环执行 - 修复了参数不匹配、重复释放连接、订阅路径错误等细节问题
内容的提问来源于stack exchange,提问作者Sahil Sahni
相关产品推荐
相关产品推荐

