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

SocketIO连接事件调用Celery任务超时问题求助

问题分析与修复方案

核心问题

  1. 同步调用Celery任务导致连接超时:直接调用celeryTask(1,2)是在当前线程同步执行函数,任务内的while True循环会无限阻塞connect事件处理流程,WebSocket客户端因长时间未收到连接确认而超时。
  2. 不合理的SocketIO实例创建与阻塞操作:Celery任务内重复创建SocketIO实例,且使用time.sleep(60)会阻塞Celery Worker线程,影响任务并发能力。

修复后的代码

from flask import Flask
from flask_socketio import SocketIO
import eventlet
from celery import Celery
import time

eventlet.monkey_patch(socket=True)
app = Flask(__name__)
app.config['SECRET_KEY'] = 'secret'
# 初始化SocketIO,使用Redis作为消息队列实现跨进程通信
socketio = SocketIO(app, async_mode='eventlet', logger=True, engineio_logger=True, message_queue='redis://127.0.0.1:6379')

celery = Celery(app.name, broker='redis://127.0.0.1:6379')
celery.conf.update(app.config)

@app.route('/')
def home():
    return 'Hello World!'

@socketio.on('connect')
def connect():
    print('Client connected, calling celery task...')
    # 使用delay()异步提交Celery任务,避免阻塞connect事件
    celeryTask.delay(1, 2)

@celery.task()
def celeryTask(x, y):
    print('Celery task called!')
    # 创建仅用于发送消息的SocketIO实例(依赖消息队列)
    task_socketio = SocketIO(message_queue='redis://127.0.0.1:6379')
    try:
        while True:
            # 通过消息队列发送事件,主进程的SocketIO会转发给客户端
            task_socketio.emit('add', {'data': x + y})
            # 使用eventlet.sleep避免阻塞Worker线程
            eventlet.sleep(60)
    finally:
        # 关闭SocketIO连接释放资源
        task_socketio.disconnect()

if __name__ == '__main__':
    socketio.run(app, debug=True)

关键修改点

  • 异步提交Celery任务:用celeryTask.delay(1,2)替代直接调用,将任务提交到Celery队列,connect事件可快速完成,客户端能正常建立连接。
  • 优化SocketIO实例:Celery任务内创建仅用于发送消息的SocketIO实例(无需绑定Flask app),通过Redis消息队列与主进程的SocketIO通信,实现跨进程事件推送。
  • 替换阻塞睡眠:用eventlet.sleep(60)替代time.sleep(60),避免阻塞Celery Worker的协程,保证Worker能处理其他任务。

运行注意事项

  1. 启动Redis服务,确保消息队列可用。
  2. 启动Celery Worker(需指定eventlet池):
celery -A app.celery worker -l info -P eventlet
  1. 启动Flask-SocketIO应用:
python app.py

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 09:31:00