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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 19:23:15