Asyncio信号量与队列:多消费者队列共享及信号量限制失效问题
Asyncio多消费者信号量不生效的原因及解决办法
嘿,我来帮你捋捋这个问题~你说consumer2的1个条目限制没生效,大概率是信号量的作用域不对,或者没有正确包裹任务处理逻辑,我给你拆解下原因和修复方案:
核心问题分析
你要给每个消费者类型设置专属的并发限制,得满足两个关键条件:
- 同类型的消费者必须共享同一个信号量实例——如果每个consumer2任务自己创建
Semaphore(1),那每个任务都能同时处理1个,等于没做全局限制; - 必须把实际处理任务的代码段用信号量的上下文管理器包裹——要是只在函数开头获取一次信号量,那整个消费者会一直占着信号量,其他同类型任务根本没法执行。
修正后的完整代码
下面是调整后的示例代码,你可以对照自己的代码看差异:
import asyncio import random filenames = {"A":1,"B":2,"C":1,"D":2,"E":1,"F":1,"G":2,"H":1,"I":2,"J":2} async def producer(f, d, q): print(f"producing {f}") # 模拟生产耗时 await asyncio.sleep(random.randint(0, 1)) await q.put((f, d)) async def consumer1(q, sem): while True: f, d = await q.get() # 用信号量包裹处理逻辑,确保同一时间最多3个consumer1任务在处理 async with sem: print(f"Consumer 1 processing {f} (type {d})") await asyncio.sleep(random.randint(1, 3)) print(f"Consumer 1 finished {f}") q.task_done() async def consumer2(q, sem): while True: f, d = await q.get() # 这里就是关键:把处理逻辑放在信号量上下文里,确保同一时间只有1个consumer2任务在处理 async with sem: print(f"Consumer 2 processing {f} (type {d})") await asyncio.sleep(random.randint(1, 3)) print(f"Consumer 2 finished {f}") q.task_done() async def main(): q = asyncio.Queue() # 为不同消费者类型创建专属信号量:consumer1允许同时3个,consumer2仅允许1个 sem_c1 = asyncio.Semaphore(3) sem_c2 = asyncio.Semaphore(1) # 启动所有生产者任务 producers = [asyncio.create_task(producer(f, d, q)) for f, d in filenames.items()] # 启动消费者任务:比如2个consumer1、3个consumer2,都受各自信号量限制 consumers = [] for _ in range(2): consumers.append(asyncio.create_task(consumer1(q, sem_c1))) for _ in range(3): consumers.append(asyncio.create_task(consumer2(q, sem_c2))) # 等待所有生产者完成生产 await asyncio.gather(*producers) # 等待队列中所有任务处理完毕 await q.join() # 取消无限循环的消费者任务 for task in consumers: task.cancel() await asyncio.gather(*consumers, return_exceptions=True) if __name__ == "__main__": asyncio.run(main())
关键细节说明
- 信号量的共享方式:在
main函数里创建信号量,然后传给对应的消费者任务,这样同类型的所有消费者都会共用同一个信号量,才能实现全局的并发限制; - 信号量的正确使用:用
async with sem:包裹实际处理任务的代码段,这样每次处理一个任务时都会自动获取信号量,处理完自动释放——如果把信号量的acquire放在循环外面,那整个消费者任务会一直占用信号量,其他同类型任务根本无法执行; - 队列的task_done:处理完每个任务后一定要调用
q.task_done(),否则q.join()会一直阻塞,程序无法正常结束。
内容的提问来源于stack exchange,提问作者bonzinor
相关产品推荐
相关产品推荐

