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

如何从事件中正确生成值?Python生成器消息处理方案问询

Python:将触发式消息函数转为持续生成器

你现在的场景是有一个被消息触发调用的on_message函数,只能通过yield返回消息,想要实现一个get_all_messages生成器持续产出这些消息,之前想到的文件/管道、套接字方案都有额外IO和数据转换的问题,最适合的方案是使用线程安全的队列(同步场景)或异步队列(异步场景),完全不需要额外的IO操作,直接在内存中传递数据。

同步场景实现(基于queue.Queue)

这种方案适用于on_message在普通线程中被调用的场景:

  1. 先创建一个全局或可共享的队列实例,用来暂存消息:
import queue
from typing import Any

# 线程安全的消息队列,可设置maxsize限制队列容量(0表示无限制)
msg_queue = queue.Queue(maxsize=0)
  1. 修改on_message函数,将原本yield的消息放入队列(触发on_message的调用方通常不会消费它的生成器,直接存队列更可靠):
def on_message(message: Any):
    # 消息放入队列,队列满时会阻塞等待
    msg_queue.put(message)
  1. 实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 23:31:12