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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 16:45:03