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

如何并行执行两类异步任务,避免连接阶段阻塞事件循环?

异步客户端连接与任务执行的分离优化

问题场景与原实现问题

程序需要完成两类核心任务:

  • 批量连接客户端列表
  • 利用已连接的客户端执行一系列任务

原实现代码如下:

import asyncio

async def task(worker_id):
    print(f'worker {worker_id} start working ...')
    
    # 实际客户端来自Pyrogram库
    client = Client(worker_id)

    await client.connect()

    for job in jobs:
        ... # 执行任务的await代码
        print(f'job {job.id} finished!')
        

async def main():
    tasks = [task(_) for _ in range(10)]
    await asyncio.gather(*tasks)

asyncio.run(main())

这段代码的问题在于:虽然用asyncio.gather并发启动了10个task协程,但每个task里的await client.connect()会让协程挂起,事件循环会依次处理这些连接操作,导致所有客户端完成连接后,才会开始执行后续的任务逻辑,连接阶段阻塞了任务执行的启动。

关于asyncio.Queue方案的可行性分析

你给出的示例代码思路方向是对的,但存在几个明显问题:

  1. task2函数定义了多余的client参数,实际逻辑是从队列获取客户端
  2. 只启动了一个task2协程,只能处理一个客户端,无法利用多客户端并发执行任务
  3. task函数是串行连接客户端,没有并发处理连接,效率低下
  4. 缺少队列结束标记,工作协程会陷入无限循环无法正常退出

优化后的实现方案

基于asyncio.Queue可以完美实现连接与任务执行的解耦,核心思路是:

  • 并发启动多个连接协程,每个客户端连接完成后立即放入队列
  • 启动多个工作协程,持续从队列中取已连接的客户端执行任务
  • 所有连接完成后,向队列放入结束标记,让工作协程可以正常退出

优化后的代码示例:

import asyncio
# 假设已导入Pyrogram的Client类,且定义了jobs任务列表、clients客户端ID列表

async def connect_worker(client_id, queue):
    print(f'worker {client_id} 开始连接...')
    client = Client(client_id)
    await client.connect()
    print(f'worker {client_id} 连接完成')
    await queue.put(client)

async def task_worker(queue):
    while True:
        client = await queue.get()
        # 收到结束标记则退出循环
        if client is None:
            queue.task_done()
            break
        try:
            print(f'开始使用客户端{client.id}执行任务')
            for job in jobs:
                await client.do_job(job)
                print(f'任务 {job.id} 完成!')
        finally:
            # 标记队列任务完成,避免队列阻塞
            queue.task_done()
            # 可选:任务完成后关闭客户端,释放资源
            await client.disconnect()

async def main():
    queue = asyncio.Queue()
    client_count = 10  # 或使用len(clients)获取实际客户端数量
    # 启动所有连接协程
    connect_tasks = [connect_worker(i, queue) for i in range(client_count)]
    # 启动与客户端数量匹配的任务工作协程,最大化并发效率
    task_workers = [asyncio.create_task(task_worker(queue)) for _ in range(client_count)]
    
    # 等待所有客户端连接完成
    await asyncio.gather(*connect_tasks)
    # 向队列放入对应数量的结束标记,每个工作协程对应一个
    for _ in range(client_count):
        await queue.put(None)
    # 等待所有任务工作协程执行完成
    await asyncio.gather(*task_workers)

asyncio.run(main())

关键优化点说明:

  • 并发连接:用asyncio.gather并发启动所有连接协程,每个客户端连接完成后立刻进入队列,不会互相阻塞
  • 多工作协程:启动多个任务工作协程,同时处理队列中的客户端,充分利用异步并发能力
  • 优雅退出:通过向队列放入None作为结束标记,让工作协程可以正常退出,避免无限循环
  • 资源清理:在任务执行完成后关闭客户端,避免资源泄漏

内容的提问来源于stack exchange,提问作者Hamidreza

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 09:09:55