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

基于FastAPI+WebSocket实现Pyomo求解器日志实时流式输出的问题

解决方案:实时推送Pyomo求解器日志到WebSocket

要实现求解器日志的近乎实时推送,核心是实时捕获求解器输出并通过WebSocket持续发送给前端,以下分两种场景给出具体实现方案:


方案一:原生FastAPI + WebSocket + 子进程(轻量部署)

直接通过子进程运行求解脚本,实时捕获标准输出/错误,逐行推送给WebSocket客户端,无需额外中间件。

1. 实现步骤

  • 定义WebSocket端点,建立客户端连接
  • 用subprocess启动求解脚本,开启行缓冲模式确保实时读取日志
  • 逐行读取子进程输出,通过WebSocket推送给前端
  • 处理客户端断开、求解出错等异常场景,及时终止子进程

2. 代码示例

主服务文件 main.py

from fastapi import FastAPI, WebSocket, WebSocketDisconnect
import subprocess
import asyncio

app = FastAPI()

@app.websocket("/ws/solver_log")
async def websocket_logger(websocket: WebSocket):
    await websocket.accept()
    # 启动求解脚本,捕获 stdout/stderr
    proc = subprocess.Popen(
        ["python", "solver_script.py"],
        stdout=subprocess.PIPE,
        stderr=subprocess.STDOUT,
        text=True,
        bufsize=1,  # 行缓冲,确保实时读取
        universal_newlines=True
    )

    try:
        # 逐行读取输出并推送
        for line in iter(proc.stdout.readline, ''):
            if line.strip():
                await websocket.send_text(line.strip())
                await asyncio.sleep(0.05)  # 控制推送频率,避免前端压力过大
        
        # 求解完成后发送结束信号
        proc.wait()
        await websocket.send_text(f"✅ 求解完成,返回码: {proc.returncode}")
    except WebSocketDisconnect:
        # 客户端断开,终止子进程
        proc.terminate()
        proc.wait()
    except Exception as e:
        await websocket.send_text(f"❌ 出错: {str(e)}")
        proc.terminate()
        proc.wait()

求解脚本 solver_script.py

from pyomo.environ import ConcreteModel, Var, Objective, Constraint, SolverFactory

# 构建示例模型(实际替换为你的业务模型)
model = ConcreteModel()
model.x = Var(within=NonNegativeReals)
model.y = Var(within=NonNegativeReals)
model.obj = Objective(expr=2*model.x + 3*model.y)
model.con1 = Constraint(expr=3*model.x + 4*model.y >= 1)
model.con2 = Constraint(expr=model.x + 2*model.y >= 1)

# 开启tee=True,让求解器日志输出到stdout
solver = SolverFactory('glpk')
result = solver.solve(model, tee=True)

方案二:FastAPI + Celery + Redis + WebSocket(高并发场景)

当需要处理大量求解任务时,用Celery做异步任务队列,Redis做消息中间件,通过Redis发布订阅实现日志的实时推送。

1. 实现步骤

  • 配置Celery和Redis,搭建异步任务环境
  • Celery任务中捕获求解器日志,推送到Redis指定频道
  • FastAPI的WebSocket端点订阅Redis频道,收到日志后推送给前端
  • 前端通过任务ID订阅对应日志流

2. 代码示例

Celery配置文件 celery_config.py

from celery import Celery
import redis

app = Celery(
    'solver_tasks',
    broker='redis://localhost:6379/0',
    backend='redis://localhost:6379/0'
)
redis_client = redis.Redis(host='localhost', port=6379, db=0)

Celery任务文件 tasks.py

from celery_config import app, redis_client
from pyomo.environ import ConcreteModel, Var, Objective, Constraint, SolverFactory
import sys
from io import StringIO

@app.task(bind=True)
def run_solver(self, task_id):
    # 重定向stdout捕获日志
    old_stdout = sys.stdout
    sys.stdout = log_buffer = StringIO()

    try:
        # 构建模型(支持传入参数动态生成)
        model = ConcreteModel()
        model.x = Var(within=NonNegativeReals)
        model.y = Var(within=NonNegativeReals)
        model.obj = Objective(expr=2*model.x + 3*model.y)
        model.con1 = Constraint(expr=3*model.x + 4*model.y >= 1)
        model.con2 = Constraint(expr=model.x + 2*model.y >= 1)

        solver = SolverFactory('glpk')
        solver.solve(model, tee=True)

        # 逐行推送日志到Redis频道
        log_buffer.seek(0)
        for line in log_buffer.readlines():
            if line.strip():
                redis_client.publish(f"solver_log_{task_id}", line.strip())
        
        redis_client.publish(f"solver_log_{task_id}", "✅ 求解完成")
    except Exception as e:
        redis_client.publish(f"solver_log_{task_id}", f"❌ 出错: {str(e)}")
    finally:
        sys.stdout = old_stdout

FastAPI主服务 main.py

from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from celery_config import app as celery_app, redis_client
import asyncio
import uuid

app = FastAPI()

@app.post("/start_solver")
async def start_solver_task():
    task_id = str(uuid.uuid4())
    celery_app.send_task('tasks.run_solver', args=[task_id])
    return {"task_id": task_id}

@app.websocket("/ws/solver_log/{task_id}")
async def websocket_logger(websocket: WebSocket, task_id: str):
    await websocket.accept()
    # 订阅对应任务的Redis日志频道
    pubsub = redis_client.pubsub()
    pubsub.subscribe(f"solver_log_{task_id}")

    async def listen_redis():
        while True:
            msg = pubsub.get_message(ignore_subscribe_messages=True)
            if msg:
                await websocket.send_text(msg['data'].decode('utf-8'))
            await asyncio.sleep(0.1)

    try:
        await listen_redis()
    except WebSocketDisconnect:
        pubsub.unsubscribe(f"solver_log_{task_id}")
    except Exception as e:
        await websocket.send_text(f"❌ 连接出错: {str(e)}")
        pubsub.unsubscribe(f"solver_log_{task_id}")

前端示例(通用)

简单HTML页面,连接WebSocket实时展示日志:

<!DOCTYPE html>
<html>
<head>
    <title>求解器日志实时查看</title>
</head>
<body>
    <h1>求解器日志</h1>
    <div id="log-container" style="height: 400px; overflow-y: auto; border: 1px solid #ccc; padding: 10px;"></div>
    <button onclick="connectLog()">连接日志</button>
    <script>
        let ws;
        // 从后端/输入框获取实际task_id(方案二用)
        const taskId = "替换为你的任务ID";

        function connectLog() {
            // 方案一用 ws://localhost:8000/ws/solver_log
            // 方案二用 ws://localhost:8000/ws/solver_log/${taskId}
            ws = new WebSocket(`ws://localhost:8000/ws/solver_log/${taskId}`);
            
            ws.onmessage = (event) => {
                const logDiv = document.getElementById('log-container');
                logDiv.innerHTML += `<p>${event.data}</p>`;
                logDiv.scrollTop = logDiv.scrollHeight; // 自动滚动到底部
            };

            ws.onclose = () => console.log("日志连接已关闭");
            ws.onerror = (err) => console.error("连接出错:", err);
        }
    </script>
</body>
</html>

关键优化点

  • 日志捕获方式:优先用求解器的tee=True输出到stdout,部分商业求解器(如Gurobi/CPLEX)支持日志回调函数,可更精准控制日志输出
  • 缓冲设置:子进程必须开启行缓冲(bufsize=1),避免日志被批量缓存无法实时推送
  • 资源隔离:高并发场景下用Celery+Redis实现任务解耦,避免求解任务阻塞FastAPI主进程
  • 异常处理:客户端断开连接时需及时终止求解任务,避免资源浪费

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 06:04:56