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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 15:57:04