如何在Python异步SocketIO多用户应用中实现进程并行执行
解决异步SocketIO服务器中用户进程串行执行的问题
问题根源
你当前的代码中,trainmodel是同步阻塞函数,在异步的trainit处理函数里直接调用会卡住整个异步事件循环——事件循环必须等待当前训练任务完成后,才能处理后续的Socket请求,因此所有用户的进程只能串行执行。
解决方案
将同步的训练任务放到线程池/进程池中执行,让异步事件循环可以同时处理多个用户请求。推荐使用asyncio.to_thread(Python 3.9+)或loop.run_in_executor,将阻塞任务从事件循环中剥离到后台线程/进程。
修改后的完整代码
import asyncio import socketio from aiohttp import web # 导入你的训练依赖(比如transformers相关模块) sio = socketio.AsyncServer(cors_allowed_origins="*") app = web.Application() sio.attach(app) def trainmodel(x, y, username, code): # 耗时的模型训练逻辑(约10秒) args = TrainingArguments( output_dir=f"output/{username}", # 按用户划分输出目录,避免文件冲突 num_train_epochs=10, per_device_train_batch_size=10, optim="adamw_torch" ) # 关键:确保model、train_dataset是当前用户的独立实例,不要用全局变量,避免多线程资源冲突 trainer = Trainer( model=model, args=args, train_dataset=train_dataset, compute_metrics=compute_metrics ) trainer.train() # 训练完成后返回结果(可选) return trainer.evaluate() @sio.on("sendbackends") async def trainit(sid, data, code): # 从请求数据中提取用户参数 x = data.get('x') y = data.get('y') username = data.get('username') # 使用asyncio.to_thread将同步任务丢到后台线程执行,不阻塞事件循环 # Python 3.9以下版本用loop.run_in_executor替代(见下方说明) training_result = await asyncio.to_thread(trainmodel, x, y, username, code) # 训练完成后给当前用户发送结果通知 await sio.emit("training_done", { "username": username, "result": training_result }, room=sid) @sio.on("bottoken") async def handle(sid): await sio.emit("redirect", namespace="/", room=sid) if __name__ == '__main__': web.run_app(app, port=5000)
适配Python 3.8及以下版本
如果你的Python版本低于3.9,用loop.run_in_executor替代asyncio.to_thread:
@sio.on("sendbackends") async def trainit(sid, data, code): x = data.get('x') y = data.get('y') username = data.get('username') loop = asyncio.get_running_loop() # None表示使用默认线程池,也可以自定义ThreadPoolExecutor training_result = await loop.run_in_executor(None, trainmodel, x, y, username, code) await sio.emit("training_done", { "username": username, "result": training_result }, room=sid)
CPU密集型任务优化
如果训练是纯CPU密集型任务,线程受GIL限制无法充分利用多核,可以改用进程池:
from concurrent.futures import ProcessPoolExecutor # 初始化进程池,大小建议等于CPU核心数 executor = ProcessPoolExecutor(max_workers=4) @sio.on("sendbackends") async def trainit(sid, data, code): x = data.get('x') y = data.get('y') username = data.get('username') loop = asyncio.get_running_loop() training_result = await loop.run_in_executor(executor, trainmodel, x, y, username, code) await sio.emit("training_done", { "username": username, "result": training_result }, room=sid)
注意:进程池要求传递给trainmodel的参数必须是可序列化的(支持pickle)。
关键注意事项
- 资源隔离:每个用户的模型、数据集必须是独立实例,禁止使用全局共享对象,否则会出现数据竞争或训练结果错误。
- 资源限制:根据服务器硬件配置调整线程/进程池大小,避免因任务过多导致内存或CPU耗尽。
- 状态追踪:如果需要给用户展示训练进度,需额外实现进度上报逻辑(比如在训练过程中定时通过Socket发送进度数据)。
内容的提问来源于stack exchange,提问作者tidekis doritos
相关产品推荐
相关产品推荐

