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

如何在Python SocketIO异步服务器中并行执行用户进程?

解决Python SocketIO异步服务器进程串行阻塞问题

问题原因

当前代码中trainmodel是同步阻塞函数,在异步的trainit处理函数里直接调用会卡住整个事件循环,导致后续的Socket请求必须等待当前训练任务完成才能被处理,最终表现为所有用户进程串行执行。

解决方案

通过将同步阻塞的训练任务放到线程池或进程池中执行,让事件循环可以继续处理其他请求,实现并行运行。

方法一:使用线程池(适合IO密集型任务)

利用asyncio.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):
    # 预处理数据、加载模型等操作
    # ...
    args = TrainingArguments(
        output_dir=f"output_{username}",  # 按用户区分输出目录,避免文件冲突
        num_train_epochs=10,
        per_device_train_batch_size=10,
        optim="adamw_torch"
    )
    trainer = Trainer(
        model=model,
        args=args,   
        train_dataset=train_dataset,
        compute_metrics=compute_metrics
    )
    trainer.train()
    return f"用户{username}的模型训练完成"

@sio.on("sendbackends") 
async def trainit(sid, data, code):
    # 从请求数据中解析所需参数(根据实际业务调整)
    x = data.get('x')
    y = data.get('y')
    username = data.get('username')
    
    # 将同步训练任务放到线程池执行,不阻塞事件循环
    train_result = await asyncio.get_event_loop().run_in_executor(
        None, trainmodel, x, y, username, code
    )
    
    # 训练完成后向当前客户端发送结果通知
    await sio.emit("train_finished", {"result": train_result}, to=sid)

@sio.on("bottoken")
async def handle():
    await sio.emit("redirect", namespace="/")

if __name__ == '__main__': 
    web.run_app(app, port=5000)

方法二:使用进程池(适合CPU密集型的AI模型训练)

对于CPU密集型的训练任务,Python的GIL会限制线程的并行效率,此时改用进程池更合适。

修改要点:

  1. 导入ProcessPoolExecutor创建进程池
  2. 将训练任务委托给进程池执行

示例代码片段:

from concurrent.futures import ProcessPoolExecutor

# 初始化进程池,根据CPU核心数设置最大进程数
executor = ProcessPoolExecutor(max_workers=4)

# 在trainit函数中替换为进程池执行
@sio.on("sendbackends") 
async def trainit(sid, data, code):
    x = data.get('x')
    y = data.get('y')
    username = data.get('username')
    
    train_result = await asyncio.get_event_loop().run_in_executor(
        executor, trainmodel, x, y, username, code
    )
    
    await sio.emit("train_finished", {"result": train_result}, to=sid)

注意事项

  • 使用进程池时,trainmodel函数及其内部的变量必须支持序列化(可被pickle处理),全局模型建议在每个进程内独立加载,避免进程间资源冲突。
  • 调整max_workers的数量,避免过多进程占用系统资源导致性能下降。
  • 如果训练需要GPU,需确保每个进程能正确访问GPU设备(如通过设置CUDA_VISIBLE_DEVICES分配不同GPU)。

内容的提问来源于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.26 14:25:33