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

Flask-SocketIO+eventlet猴子补丁场景下Celery任务无法运行求助

更新2

解决方案在于猴子补丁的实际实现方式,详见下方我的回答。

更新1

问题出在eventlet的猴子补丁上。我对猴子补丁了解不深,无法完全理解原因。不打猴子补丁时Celery任务可以运行,但此时无法在其他进程中使用SocketIO实例,因为不打补丁就无法将Redis作为SocketIO消息队列。尝试使用gevent仍有问题,后续会更新结果。

另外需要说明,我必须将Celery对象实例化方式改为Celery("backend"),而非Celery(app.import_name)、Celery(app.name)或Celery(__name__),才能让未打猴子补丁的任务运行。由于我的任务不使用应用上下文的任何内容,实际上甚至不需要make_celery.py模块,可直接在backend.py中实例化。

我还尝试了Redis中的不同数据库,以为可能存在冲突,但这并非问题所在。

我还通过telnet对Celery任务进行远程调试。再次验证:不打猴子补丁时任务可运行,外部socketio对象存在,但无法与主服务器通信以发送消息;打猴子补丁时任务根本无法运行。

使用gevent并打猴子补丁时,应用甚至无法启动。

目标

运行一个实时Web应用,通过Socket.IO将并行进程生成的数据推送给客户端。我使用Flask和Flask-SocketIO。

已尝试方案

最初我按照文档描述从外部进程发送消息(原始最小可运行示例见对应仓库),但存在bug。具体来说,当数据流对象在if __name__ == "__main__:"块中实例化并启动外部进程时,运行完美;但在Socket.IO事件中按需实例化并启动时则失败。大量研究指向一个开放的eventlet问题,表明eventlet与多进程兼容性不佳。之后我尝试使用gevent,短期内可行,但长时间运行(如12小时)仍存在bug。

某个回答引导我在应用中尝试Celery,但自此陷入困境。具体问题是,任务状态会显示为pending一段时间(推测是默认超时时间),然后显示失败。以debug日志级别运行worker时,错误信息为Received unregistered task of type '__main__.stream_data'.。我尝试了所有能想到的启动worker和注册任务的方式。我感到困惑,因为Celery实例与任务定义在同一作用域中,且按照无数教程和示例的做法,我使用celery -A backend.cel worker -l DEBUG命令启动worker,指定从backend.py模块中的Celery实例获取任务(至少我是这么理解这个命令的)。

当前项目状态
.                                                                                                                                                                                                                 
├── backend.py                                                                                                                                                                                                     
├── static                                                                                                                                                                                                         
│   └── js                                                                                                                                                                                                         
│       ├── main.js                                                                                                                                                                                                
│       ├── socket.io.min.js                                                                                                                                                                                       
│       └── socket.io.min.js.map                                                                                                                                                                                   
└── templates                                                                                                                                                                                                         
    └── index.html

backend.py

import eventlet
eventlet.monkey_patch()
# ^^^ 注释/取消注释以控制任务是否运行

from random import randrange
import time

from redis import Redis

from flask import Flask, render_template, request
from flask_socketio import SocketIO
from celery import Celery
from celery.contrib import rdb


def message_queue(db):
    # 曾以为可能是Redis数据库冲突,因此允许选择数据库
    # 但这并非问题所在
    return f"redis://localhost:6379/{db}"

app = Flask(__name__)
socketio = SocketIO(app, message_queue=message_queue(0))

cel = Celery("backend", broker=message_queue(0), backend=message_queue(0))

@app.route('/')
def index():
    return render_template("index.html")

@socketio.on("start_data_stream")
def start_data_stream():
    socketio.emit("new_data", {"value" :  666}) # <<<  sanity check,socket服务器在此处正常工作
    stream_data.delay(request.sid)

@cel.task()
def stream_data(sid):

    data_socketio = SocketIO(message_queue=message_queue(0))
    i = 1

    while i <= 100:
        value = randrange(0, 10000, 1) / 100
        data_socketio.emit("new_data", {"value" :  value})
        i += 1
        time.sleep(0.01)
    
    # rdb.set_trace() # <<<< 调试时可注释/取消注释,参考Celery调试文档

    return i, value


if __name__ == "__main__":

    r = Redis()
    r.flushall()

    if r.ping():
        pass
    else:
        raise Exception("需要安装Redis,请检查redis-server.service是否运行!")

    ip = "192.168.1.8" # 在此处填入局域网地址
    port = 8080

    socketio.run(app, host=ip, port=port, use_reloader=False, debug=True)

index.html

<!DOCTYPE html>
<html>
<head>

  <title>Minimal Example</title>
  <script src="{{ url_for('static', filename='js/socket.io.min.js') }}"></script>

</head>
<body>

    <button id="start" onclick="button_handler()">Start Stream</button>
    <span id="data"></span>

    <script type="text/javascript" src="{{ url_for('static', filename='js/main.js') }}"></script>

</body>
</html>

main.js

var socket =  io(location.origin);
var span = document.getElementById("data");

function button_handler() {
    socket.emit("start_data_stream");
}

socket.on("new_data", function(data) {
    span.innerHTML = data.value;
});

依赖包

Package          Version
---------------- -------
amqp             5.1.1
async-timeout    4.0.2
bidict           0.22.0
billiard         3.6.4.0
celery           5.2.7
click            8.1.3
click-didyoumean 0.3.0
click-plugins    1.1.1
click-repl       0.2.0
Deprecated       1.2.13
dnspython        2.2.1
eventlet         0.33.2
Flask            2.2.2
Flask-SocketIO   5.3.2
gevent           22.10.2
gevent-websocket 0.10.1
greenlet         2.0.1
itsdangerous     2.1.2
Jinja2           3.1.2
kombu            5.2.4
MarkupSafe       2.1.1
packaging        22.0
pip              22.3.1
prompt-toolkit   3.0.36
python-engineio  4.3.4
python-socketio  5.7.2
pytz             2022.6
redis            4.4.0
setuptools       58.1.0
six              1.16.0
vine             5.0.0
wcwidth          0.2.5
Werkzeug         2.2.2
wrapt            1.14.1
zope.event       4.6
zope.interface   5.5.2
问题

这仍然是eventlet的问题吗?某个回答让我认为Celery是eventlet问题的解决方案,且Celery甚至不需要eventlet就能工作。但eventlet似乎在我的代码中根深蒂固,因为不打猴子补丁的话,Redis相关功能完全无法工作。另外,Flask-SocketIO文档表明它会自动寻找eventlet,那么仅仅在Celery任务中实例化外部SocketIO服务器就会引入eventlet,对吗?我是否还有其他操作错误?或许有更好的worker和任务调试方法?

任何帮助都将不胜感激,谢谢!


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 19:45:31