基于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
相关产品推荐
相关产品推荐

