如何将Python多线程生产者消费者模型改造为async/await异步模式
实现方案
核心改造思路
- 保留原有线程生产者逻辑(适合处理网络请求这类阻塞IO场景),新增asyncio安全的消息队列作为对外交互层
- 封装所有线程管理、消息转发逻辑到独立类,对外仅暴露
async consumer_receives_a_message()接口和启停方法 - 跨线程投递消息时用asyncio提供的线程安全方法,避免事件循环冲突
完整封装代码
import threading import random import logging import asyncio import concurrent.futures class MessagePipeline: def __init__(self, max_queue_size: int = 10): self._stop_event = threading.Event() self._async_queue = asyncio.Queue(maxsize=max_queue_size) self._executor = concurrent.futures.ThreadPoolExecutor(max_workers=1) self._producer_task = None def _producer(self): """原有生产者逻辑,保留在线程中运行""" while not self._stop_event.is_set(): message = random.randint(1, 101) logging.info("Producer got message: %s", message) # 线程安全地投递消息到async队列 asyncio.run_coroutine_threadsafe( self._async_queue.put(message), asyncio.get_event_loop() ) logging.info("Producer received event. Exiting") def start(self): """启动生产者线程""" self._producer_task = self._executor.submit(self._producer) async def consumer_receives_a_message(self): """对外暴露的async消费接口""" return await self._async_queue.get() async def stop(self): """停止生产者,清空剩余消息""" self._stop_event.set() # 等待生产者线程退出 await asyncio.wrap_future(self._producer_task) # 清空队列剩余消息(可选,按业务需求调整) while not self._async_queue.empty(): self._async_queue.get_nowait() self._async_queue.task_done() await self._async_queue.join() self._executor.shutdown() logging.info("Pipeline stopped successfully")
使用示例
完全符合要求的调用形式:
async def aNewFunction(pipeline: MessagePipeline): message = await pipeline.consumer_receives_a_message() # 调用方自定义消息处理逻辑 do_somthing_to_process_message(message) # 主入口示例 async def main(): format = "%(asctime)s: %(message)s" logging.basicConfig(format=format, level=logging.INFO, datefmt="%H:%M:%S") pipeline = MessagePipeline(max_queue_size=10) pipeline.start() try: # 可同时运行多个消费协程,也可以按业务逻辑循环调用aNewFunction while True: await aNewFunction(pipeline) # 其他业务逻辑 except asyncio.CancelledError: logging.info("Main: about to stop pipeline") await pipeline.stop() if __name__ == "__main__": asyncio.run(main())
原有代码问题说明
你贴出的初始消费者逻辑存在bug:do_somthing_to_process_message(message) 写在while循环外,仅会处理队列的最后一条消息,改造后的方案去掉了硬编码的消费逻辑,完全交由调用方自定义处理。
内容的提问来源于stack exchange,提问作者N Johnson
相关产品推荐
相关产品推荐

