Flask-SocketIO+eventlet猴子补丁场景下Celery任务无法运行求助
解决方案在于猴子补丁的实际实现方式,详见下方我的回答。
问题出在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

