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

如何通过Monkeypatch将LOGGER输出转为yield实现FastAPI流式传输?

问题

我使用一个主要从CLI调用的库,其中有个长运行函数通过LOGGER.info输出运行状态,示例代码如下:

def long_running_function():
    ...
    LOGGER.info("first thing done")
    ...
    LOGGER.info("second thing done")
    ...

现在想从网页触发该函数,并将这些日志流式传输到浏览器。我了解到FastAPI中应该使用Server-Sent Events(SSE),但SSE要求函数通过yield返回消息,示例代码如下:

async def long_running_stream():
    while True:
        yield "data: message\n\n"

请问是否可以通过Monkeypatch将这些LOGGER语句转为yield语句?在long_running_stream中调用long_running_function是否可行?或者有没有更好的实现方案?

解决方案

方案一:猴子补丁快速实现日志转SSE

可以通过替换日志记录器的info方法,将日志内容转为SSE格式消息并缓存,再在SSE接口中逐行yield输出。需要注意同步长运行函数必须放在后台线程执行,避免阻塞FastAPI事件循环。

示例代码:

import logging
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import threading
from queue import Queue

app = FastAPI()
LOGGER = logging.getLogger(__name__)

# 用于缓存日志消息的队列
log_queue = Queue()

def patched_info(message, *args, **kwargs):
    # 转换为SSE格式消息
    sse_msg = f"data: {message % args}\n\n"
    log_queue.put(sse_msg)
    # 保留原有日志输出(可选)
    original_info(message, *args, **kwargs)

# 保存原始的info方法
original_info = LOGGER.info

def long_running_function():
    # 模拟长运行任务
    import time
    LOGGER.info("first thing done")
    time.sleep(2)
    LOGGER.info("second thing done")
    time.sleep(2)
    LOGGER.info("task completed")

async def log_streamer():
    # 监听队列,有消息则yield
    while True:
        msg = log_queue.get()
        if msg is None:  # 用None标记任务结束
            break
        yield msg

@app.get("/stream-task")
async def stream_task():
    # 替换日志的info方法
    LOGGER.info = patched_info
    # 后台线程执行长运行函数,任务结束后向队列发送结束信号
    thread = threading.Thread(target=lambda: (long_running_function(), log_queue.put(None)))
    thread.start()
    # 返回SSE流响应
    return StreamingResponse(log_streamer(), media_type="text/event-stream")

方案二:自定义日志Handler(更规范)

相比猴子补丁,自定义日志Handler是更优雅的方案,无需修改原有日志方法的指向,只需给日志记录器添加专门的Handler来捕获日志并转发到SSE流,隔离性更好。

示例代码:

import logging
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import threading
from queue import Queue

app = FastAPI()
LOGGER = logging.getLogger(__name__)

class SSELogHandler(logging.Handler):
    def __init__(self, queue):
        super().__init__()
        self.queue = queue

    def emit(self, record):
        # 格式化日志消息并转为SSE格式
        formatted_msg = self.format(record)
        sse_msg = f"data: {formatted_msg}\n\n"
        self.queue.put(sse_msg)

def long_running_function():
    # 模拟长运行任务
    import time
    LOGGER.info("first thing done")
    time.sleep(2)
    LOGGER.info("second thing done")
    time.sleep(2)
    LOGGER.info("task completed")

async def log_streamer(queue):
    while True:
        msg = queue.get()
        if msg is None:
            break
        yield msg

@app.get("/stream-task")
async def stream_task():
    log_queue = Queue()
    # 创建并添加自定义Handler
    sse_handler = SSELogHandler(log_queue)
    LOGGER.addHandler(sse_handler)
    
    # 后台线程执行任务,结束后发送结束信号
    thread = threading.Thread(target=lambda: (long_running_function(), log_queue.put(None)))
    thread.start()
    
    try:
        # 返回SSE流响应
        return StreamingResponse(log_streamer(log_queue), media_type="text/event-stream")
    finally:
        # 任务结束后移除Handler,避免影响其他请求
        LOGGER.removeHandler(sse_handler)

关键注意事项

  • 同步长运行函数必须放在后台线程执行,否则会阻塞FastAPI的事件循环,导致服务无法处理其他请求。
  • 如果能将长运行函数改为异步实现,可直接在异步上下文调用,无需额外线程,效率更高。
  • 猴子补丁方式快捷但侵入性强,若其他地方使用同一LOGGER会受影响;自定义Handler的隔离性和可维护性更优。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 04:51:32