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

如何实时捕获InteractiveInterpreter执行代码的stdout输出?

实时捕获InteractiveInterpreter输出的方案

你的异步思路方向是对的,但要注意一个关键问题:InteractiveInterpreter.runcode是同步方法,里面的time.sleep这类阻塞操作会直接卡住asyncio事件循环,所以必须把代码执行放到线程池或者子进程中,才能让输出读取的协程正常运行。

可行的异步实现方案

方案1:线程池+自定义队列流(保留Interpreter环境)

这个方案适合需要多次执行代码共享同一个Interpreter环境的场景,通过自定义流把输出写入线程安全队列,再用协程实时读取:

import io
import queue
from code import InteractiveInterpreter
from contextlib import redirect_stdout
import asyncio

class RedirectStream(io.TextIOBase):
    def __init__(self, q):
        self.q = q

    def write(self, s):
        # 过滤空行,可根据需求调整
        if s.strip():
            self.q.put(s)
        return len(s)

async def realtime_output_handler(q):
    while True:
        try:
            # 非阻塞读取队列,避免卡住协程
            output = q.get(block=False)
            # 这里替换成你的gRPC实时响应逻辑
            print(f"实时输出: {output.strip()}")
            q.task_done()
        except queue.Empty:
            await asyncio.sleep(0.1)  # 降低轮询频率,减少CPU消耗

async def run_interpreter_code(code):
    output_queue = queue.Queue()
    inter = InteractiveInterpreter()

    # 启动实时输出处理协程
    output_task = asyncio.create_task(realtime_output_handler(output_queue))

    # 把同步的runcode放到线程池执行,不阻塞事件循环
    await asyncio.to_thread(
        lambda: redirect_stdout(RedirectStream(output_queue))(inter.runcode)(code)
    )

    # 等待队列中所有输出处理完毕
    await asyncio.to_thread(output_queue.join)
    # 关闭输出处理协程
    output_task.cancel()
    try:
        await output_task
    except asyncio.CancelledError:
        pass

if __name__ == "__main__":
    test_code = """import time
for i in range(3):
    print(i)
    time.sleep(1)"""
    asyncio.run(run_interpreter_code(test_code))

方案2:异步子进程+管道读取(隔离性更好)

如果不需要共享Interpreter环境,用异步子进程的方式更简洁,还能实现代码与主进程的隔离(适合执行不可信代码):

import asyncio

async def run_code_via_subprocess(code):
    # 创建异步子进程执行代码
    proc = await asyncio.create_subprocess_exec(
        "python", "-c", code,
        stdout=asyncio.subprocess.PIPE,
        stderr=asyncio.subprocess.PIPE,
        text=True
    )

    # 实时读取子进程stdout
    async for line in proc.stdout:
        # 替换为gRPC实时响应逻辑
        print(f"实时输出: {line.strip()}")

    # 等待子进程结束
    await proc.wait()

if __name__ == "__main__":
    test_code = """import time
for i in range(3):
    print(i)
    time.sleep(1)"""
    asyncio.run(run_code_via_subprocess(test_code))

其他注意事项

  • 如果你用的是gRPC异步服务(grpc.aio),上述方案可以直接集成,把输出处理逻辑替换为gRPC响应流的发送操作即可。
  • 方案1中,queue.Queue是线程安全的,所以在子线程中写入、协程中读取不会有线程安全问题。
  • 如果需要捕获stderr输出,只需在方案1中同时重定向stderr到同一个队列,或者在方案2中读取proc.stderr。

内容的提问来源于stack exchange,提问作者naivedeveloper

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 14:27:23