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

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__这类底层生成器方法。

如果需要让一个消费者接收多个生产者的数据,只需让所有生产者往同一个队列写入数据,消费者从队列读取即可;若要区分数据来源,可在传递的数据包中加入生产者标识。

二、构建网状广播拓扑(每个生产者发所有消费者)

要实现每个生产者向所有消费者广播数据,有两种简洁思路:

  1. 消费者专属队列模式:给每个消费者创建独立队列,生产者遍历所有队列写入数据,确保每个消费者都能收到完整数据。
  2. 多订阅广播队列:基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 09:41:17