Flask-SocketIO emit被Flask-SQLAlchemy查询/提交阻塞问题排查
问题描述
我在结合Flask-SocketIO和Flask-SQLAlchemy使用时遇到了阻塞问题:socketio.emit()会被数据库的查询、提交操作阻塞。运行/test路由时,数据库的查询与提交会先全部执行完成,之后所有socketio.emit()才会一次性按顺序发送,而我预期emit能与数据库操作同步进行。
刚重装了虚拟环境,使用的依赖版本如下:
- Flask-Socketio 5.2.0
- python-engineio 4.3.4
- python-socketio 5.7.2
- gevent 22.10.2
- gevent-websocket 0.10.1
Flask服务器启动时无相关警告,想确认是否配置有误,是否需要采用更并发的方式处理SQLAlchemy?
可复现测试代码
from gevent import monkey as curious_george curious_george.patch_all(thread=True) from flask_socketio import SocketIO from flask import Flask from flask_sqlalchemy import SQLAlchemy app = Flask(__name__) app.config['SQLALCHEMY_DATABASE_URI'] = "my_custom_uri" app.config['SQLALCHEMY_POOL_SIZE'] = 20 app.config['SQLALCHEMY_MAX_OVERFLOW'] = 50 db = SQLAlchemy(app) socketio = SocketIO(app) class test_table(db.Model): _id = db.Column("id",db.Integer,primary_key=True) flagged = db.Column("flagged",db.Boolean(), default=False) def __init__(self,flagged): self.flagged = flagged @app.route('/test') def test_route(): for x in range(100): test_flag = db.session.query(test_table).filter_by(_id=1).first() test_flag.flagged = not test_flag.flagged db.session.commit() # print(x) socketio.emit('x_test',x, broadcast=True) return("success",200) if __name__ == '__main__': socketio.run(app, host='0.0.0.0', port=8080)
问题分析与解决办法
核心原因
gevent基于单线程协程模型,一旦某个协程被阻塞式操作(比如同步数据库查询/提交)占据,其他任务(包括SocketIO的消息推送)都会被挂起,直到阻塞操作完成。你的代码中,数据库操作是同步阻塞的,且在主线程协程中连续执行100次,导致socketio.emit()无法及时被调度执行,只能等所有数据库操作完成后批量推送。
具体解决措施
确保使用gevent兼容的数据库驱动
不同数据库需要搭配对应的兼容驱动,确保猴子补丁能正确异步化数据库操作:- MySQL:使用
pymysql驱动,连接URI格式改为mysql+pymysql://user:pass@host/db,并安装pymysql包 - PostgreSQL:安装
psycogreen包,在补丁后添加psycogreen.gevent.patch_psycopg()来patch psycopg2驱动
驱动兼容是异步化数据库操作的基础,否则猴子补丁无法生效,数据库操作依然是阻塞的。
- MySQL:使用
将数据库操作与emit逻辑放到后台任务
利用Flask-SocketIO的后台任务功能,把循环逻辑放到后台协程中执行,让主线程立即返回响应,后台任务可以逐步执行数据库操作和emit,避免阻塞:@app.route('/test') def test_route(): def background_job(): for x in range(100): # 每次循环用上下文管理会话,自动处理提交与关闭 with db.session.begin(): test_flag = db.session.query(test_table).filter_by(_id=1).first() test_flag.flagged = not test_flag.flagged socketio.emit('x_test', x, broadcast=True) # 可选:添加小延迟,避免数据库操作过于频繁 import time time.sleep(0.1) socketio.start_background_task(background_job) return ("success", 200)验证猴子补丁生效状态
在代码开头添加验证代码,确认关键模块已被正确补丁:from gevent import monkey as curious_george curious_george.patch_all(thread=True) # 验证补丁状态 print("socket patched:", curious_george.is_module_patched('socket')) print("sqlalchemy patched:", curious_george.is_module_patched('sqlalchemy'))如果输出为
True,说明补丁生效;若为False,检查补丁是否在所有模块导入前执行(你的代码顺序是对的,但可能驱动未被覆盖)。优化SQLAlchemy连接池配置
调整连接池参数,适配协程环境:app.config['SQLALCHEMY_POOL_RECYCLE'] = 300 # 定期回收连接,避免超时 app.config['SQLALCHEMY_POOL_TIMEOUT'] = 10 # 连接超时时间这些配置能减少连接池的阻塞概率,提升协程环境下的数据库操作效率。
内容的提问来源于stack exchange,提问作者Andrew McDonald

