Python asyncio中用协程实现多生产者多消费者的规范方案
关于Python asyncio构建协程拓扑的问题
我正在学习Python asyncio,现有代码能运行,但不符合协程编程的规范。我想要构建任意处理流程拓扑(如附图所示),有两个核心疑问:
- 如何创建能接收多个协程数据的协程?是否需要使用
asend? - 如何构建每个生产者向所有消费者发送数据的网状结构?
以下是我当前实现的「每个生产者对应多个消费者」拓扑的代码:
class Producer: async def producer(self): yield "started producer" for i in range(0, 100): for connection in self.connections: await connection.asend(i) class Consumer: async def consumer(self): yield "started consumer" while True: for connection in self.producers: item = yield 0 print("item", item) async def do(): producers = [] consumers = [] for i in range(0, 10): producers.append(Producer()) for i in range(0, 10): consumers.append(Consumer()) prodawait = [] consumeawait = [] for b in consumers: cons = b.consumer() await cons.__anext__() consumeawait.append(cons) for a in producers: prod = a.producer() prodawait.append(prod) await prod.__anext__() for producer in producers: producer.connections = consumeawait for consumer in consumers: consumer.producers = prodawait while True: for consumer in consumeawait: for item in prodawait: await item.__anext__() await consumer.__anext__() import asyncio asyncio.run(do())
规范实现方案
一、接收多协程数据的协程实现
你当前用asend手动操作生成器的方式不符合asyncio的规范实践,推荐用asyncio.Queue作为协程间的通信载体——这是官方认可的生产者-消费者模型标准方案,不需要手动调用asend或__anext__这类底层生成器方法。
如果需要让一个消费者接收多个生产者的数据,只需让所有生产者往同一个队列写入数据,消费者从队列读取即可;若要区分数据来源,可在传递的数据包中加入生产者标识。
二、构建网状广播拓扑(每个生产者发所有消费者)
要实现每个生产者向所有消费者广播数据,有两种简洁思路:
- 消费者专属队列模式:给每个消费者创建独立队列,生产者遍历所有队列写入数据,确保每个消费者都能收到完整数据。
- 多订阅广播队列:基于
asyncio.Queue封装支持多订阅的队列,每个消费者订阅后均可收到生产者发送的每一条数据。
重构后的代码示例
下面是用asyncio.Queue实现的规范版本,完成每个生产者向所有消费者广播数据的逻辑:
import asyncio class Producer: def __init__(self, producer_id, consumer_queues): self.producer_id = producer_id self.consumer_queues = consumer_queues async def run(self): print(f"启动生产者 {self.producer_id}") for i in range(100): # 封装带来源标识的数据 data = {"producer_id": self.producer_id, "value": i} # 向所有消费者队列广播数据 for queue in self.consumer_queues: await queue.put(data) await asyncio.sleep(0.01) # 模拟生产间隔 class Consumer: def __init__(self, consumer_id, own_queue): self.consumer_id = consumer_id self.own_queue = own_queue async def run(self): print(f"启动消费者 {self.consumer_id}") while True: # 从自己的队列读取数据(所有生产者都会往这里写) data = await self.own_queue.get() print(f"消费者 {self.consumer_id} 收到生产者 {data['producer_id']} 的数据: {data['value']}") self.own_queue.task_done() async def main(): # 创建10个消费者的专属队列 consumer_count = 10 consumer_queues = [asyncio.Queue(maxsize=10) for _ in range(consumer_count)] # 初始化生产者:每个生产者持有所有消费者的队列 producers = [Producer(i, consumer_queues) for i in range(10)] # 初始化消费者:每个消费者持有自己的专属队列 consumers = [Consumer(i, consumer_queues[i]) for i in range(consumer_count)] # 启动所有协程任务 tasks = [] for p in producers: tasks.append(asyncio.create_task(p.run())) for c in consumers: tasks.append(asyncio.create_task(c.run())) # 等待所有任务完成(可按需求调整终止逻辑) await asyncio.gather(*tasks) if __name__ == "__main__": asyncio.run(main())
关键改进说明
- 用
asyncio.Queue替代手动生成器操作,代码更易维护、符合asyncio设计规范。 - 生产者通过遍历所有消费者队列实现广播逻辑,直观且易扩展。
- 消费者只需监听自己的队列即可接收所有生产者的广播数据,逻辑清晰。
- 避免了直接操作
__anext__、asend这类底层生成器API,这类接口通常用于框架内部实现,业务代码中应尽量避免。
内容的提问来源于stack exchange,提问作者Samuel Squire
相关产品推荐
相关产品推荐

