如何从事件中正确生成值?Python生成器消息处理方案问询
Python:将触发式消息函数转为持续生成器
你现在的场景是有一个被消息触发调用的on_message函数,只能通过yield返回消息,想要实现一个get_all_messages生成器持续产出这些消息,之前想到的文件/管道、套接字方案都有额外IO和数据转换的问题,最适合的方案是使用线程安全的队列(同步场景)或异步队列(异步场景),完全不需要额外的IO操作,直接在内存中传递数据。
同步场景实现(基于queue.Queue)
这种方案适用于on_message在普通线程中被调用的场景:
- 先创建一个全局或可共享的队列实例,用来暂存消息:
import queue from typing import Any # 线程安全的消息队列,可设置maxsize限制队列容量(0表示无限制) msg_queue = queue.Queue(maxsize=0)
- 修改
on_message函数,将原本yield的消息放入队列(触发on_message的调用方通常不会消费它的生成器,直接存队列更可靠):
def on_message(message: Any): # 消息放入队列,队列满时会阻塞等待 msg_queue.put(message)
- 实现
get_all_messages生成器,无限循环从队列取消息并yield:
def get_all_messages(): while True: # 阻塞等待新消息,拿到后立即产出 message = msg_queue.get() yield message # 标记消息处理完成(不需要等待全量消息处理时可省略) msg_queue.task_done()
异步场景实现(基于asyncio.Queue)
如果消息触发在异步环境(比如asyncio框架)中,改用异步队列:
import asyncio from typing import Any async_msg_queue = asyncio.Queue(maxsize=0) async def on_message(message: Any): await async_msg_queue.put(message) async def get_all_messages(): while True: message = await async_msg_queue.get() yield message async_msg_queue.task_done()
方案优势
- 全程内存操作,无文件/套接字的IO开销和数据转换麻烦
- 线程/协程安全,避免多环境下的消息丢失或混乱
- 可通过
maxsize控制队列容量,防止消息堆积导致内存溢出
内容的提问来源于stack exchange,提问作者vasker
相关产品推荐
相关产品推荐

