Flask-SQLAlchemy Queue Pool超限问题求助(附代码配置)
Flask应用SQLAlchemy连接池超限问题排查
错误日志
172.30.0.17 - - [16/Jun/2022 06:25:07] "GET /api/orgs HTTP/1.1" 500 - Traceback (most recent call last): File "/root/.local/lib/python3.9/site-packages/flask/app.py", line 2095, in __call__ return self.wsgi_app(environ, start_response) File "/root/.local/lib/python3.9/site-packages/flask/app.py", line 2080, in wsgi_app response = self.handle_exception(e) File "/root/.local/lib/python3.9/site-packages/flask/app.py", line 2077, in wsgi_app response = self.full_dispatch_request() File "/root/.local/lib/python3.9/site-packages/flask/app.py", line 1525, in full_dispatch_request rv = self.handle_user_exception(e) File "/root/.local/lib/python3.9/site-packages/flask/app.py", line 1523, in full_dispatch_request rv = self.dispatch_request() File "/root/.local/lib/python3.9/site-packages/flask/app.py", line 1509, in dispatch_request return self.ensure_sync(self.view_functions[rule.endpoint])(**req.view_args) File "/service/routes/get/get_org.py", line 17, in base_route query = [convert_class(item) for item in Organization.query.all()] File "/root/.local/lib/python3.9/site-packages/sqlalchemy/orm/query.py", line 2768, in all return self._iter().all() File "/root/.local/lib/python3.9/site-packages/sqlalchemy/orm/query.py", line 2903, in _iter result = self.session.execute( File "/root/.local/lib/python3.9/site-packages/sqlalchemy/orm/session.py", line 1711, in execute conn = self._connection_for_bind(bind) File "/root/.local/lib/python3.9/site-packages/sqlalchemy/orm/session.py", line 1552, in _connection_for_bind return self._transaction._connection_for_bind( File "/root/.local/lib/python3.9/site-packages/sqlalchemy/orm/session.py", line 747, in _connection_for_bind conn = bind.connect() File "/root/.local/lib/python3.9/site-packages/sqlalchemy/engine/base.py", line 3234, in connect return self._connection_cls(self, close_with_result=close_with_result) File "/root/.local/lib/python3.9/site-packages/sqlalchemy/engine/base.py", line 96, in __init__ else engine.raw_connection() File "/root/.local/lib/python3.9/site-packages/sqlalchemy/engine/base.py", line 3313, in raw_connection return self._wrap_pool_connect(self.pool.connect, _connection) File "/root/.local/lib/python3.9/site-packages/sqlalchemy/engine/base.py", line 3280, in _wrap_pool_connect return fn() File "/root/.local/lib/python3.9/site-packages/sqlalchemy/pool/base.py", line 310, in connect return _ConnectionFairy._checkout(self) File "/root/.local/lib/python3.9/site-packages/sqlalchemy/pool/base.py", line 868, in _checkout fairy = _ConnectionRecord.checkout(pool) File "/root/.local/lib/python3.9/site-packages/sqlalchemy/pool/base.py", line 476, in checkout rec = pool._do_get() File "/root/.local/lib/python3.9/site-packages/sqlalchemy/pool/impl.py", line 134, in _do_get raise exc.TimeoutError( sqlalchemy.exc.TimeoutError: QueuePool limit of size 20 overflow 20 reached, connection timed out, timeout 30.00
错误说明:连接池核心连接+溢出连接已达上限,无法获取新连接导致超时。
配置信息
class Config: # 数据库配置 SQLALCHEMY_DATABASE_URI = environ.get("DB_URI") SQLALCHEMY_ECHO = False SQLALCHEMY_TRACK_MODIFICATIONS = False SQLALCHEMY_ENGINE_OPTIONS = { "pool_size": 1, "max_overflow": 0, }
应用初始化代码
db = SQLAlchemy() # 创建Flask应用 def create_app(): app = Flask(__name__, instance_relative_config=False) # PostgreSQL数据库配置 app.config['SQLALCHEMY_DATABASE_URI'] = getenv("DB_URI") app.config.from_object('config.Config') # 配置CORS(已注释) # cors = CORS(app, resources={r"/*": {"origins": "*"}}) with app.app_context(): db.init_app(app) from service.routes.post import post_event return app app = create_app()
POST请求处理代码
@app.route('/api/orgs/<org_id>/sites/<site_id>/sessions/<session_id>/events', methods=["POST"]) def edger_post_request(org_id, site_id, session_id): print("URL参数: {}, {}, {}".format(org_id, site_id, session_id)) # 解析POST请求体 try: request_body = request.json event_data = json.loads(request_body["data"]) except: return "请求体格式错误", 400 print("请求体内容: ", request_body) # 验证请求体合法性 errors_body, has_error = validate_input(request_body) print(errors_body, has_error) if has_error: return Response(errors_body, 400) # 插入事件 event_dict = dict() try: event_dict = { "id": request_body["id"], "timestamp": request_body["timestamp"], "type": request_body["type"], "event_type": event_data["type"], "data": request_body["data"], "org_id": org_id, "site_id": site_id, "device_id": request_body["device_id"], "session_id": request_body["session_id"] } if "frame_url" in request_body.keys(): event_dict["frame_url"] = request_body["frame_url"] else: event_dict["frame_url"] = None insert_event = Event(**event_dict) db.session.add(insert_event) # 提交事件到数据库 db.session.commit() # 插入关联表数据 insert_threads(org_id, site_id, event_dict, event_data) db.session.commit() db.session.close() return {"result": "success"}, 200 except Exception as e: print("- - - - - - - - - - - - - - - - -") db.session.rollback() print("插入事件失败", e) print("- - - - - - - - - - - - - - - - -") db.session.close() return {"result": "unsuccessful"}, 409
insert_threads相关函数代码
def insert_threads(org_id, site_id, body, data): event_data_insert(body, data), event_type_insert(data["type"]), insert_org(org_id), insert_site(org_id, site_id), insert_device(org_id, site_id, body['device_id']), insert_session(org_id, site_id, body['device_id'], body['session_id'], body['timestamp']) def insert_org(org_id): org_list = getattr(g, "org_list", None) # 查询组织列表,不存在则插入 if org_list is None: check = [org.id for org in Organization.query.all()] g.org_list = check org_list = check if org_id not in org_list: insert_org = Organization(id=org_id, name="temp_name") db.session.add(insert_org) org_list.append(insert_org.id) g.org_list = org_list def insert_site(org_id, site_id): site_list = getattr(g, "site_list", None) # 查询站点列表,不存在则插入 if site_list is None: check = [site.id for site in Site.query.all()] g.site_list = check site_list = check if site_id not in site_list: insert_site = Site(id=site_id, org_id=org_id, name="temp_name") db.session.add(insert_site) site_list.append(insert_site.id) g.site_list = site_list def insert_device(org_id, site_id, device_id): device_list = getattr(g, "device_list", None) # 查询设备列表,不存在则插入 if device_list is None: check = [device.id for device in Device.query.all()] g.device_list = check device_list = check if device_id not in device_list: insert_device = Device(id=device_id, site_id=site_id, org_id=org_id) db.session.add(insert_device) device_list.append(insert_device.id) g.device_list = device_list def insert_session(org_id, site_id, device_id, session_id, timestamp): session_list = getattr(g, "session_list", None) # 查询会话列表,不存在则插入 if session_list is None: check = [session.id for session in Session.query.all()] g.session_list = check session_list = check if session_id not in session_list: insert_session = Session( id=session_id, timestamp=timestamp, org_id=org_id, site_id=site_id, device_id=device_id ) db.session.add(insert_session) session_list.append(insert_session.id) g.session_list = session_list def event_data_insert(event, event_data): # 从全局变量获取白名单,不存在则查询数据库 query_list = getattr(g, "query_list", None) if query_list is None: check = [convert_class(item) for item in Event_data_query_list.query.all()] g.query_list = check query_list = check # 遍历白名单,插入键值对到event_data表 for key in query_list: print("当前键: ", key) key_path = key["key_name"].split(".") # 检查数据对象中是否存在键,逐层取值 if key_path[0] in event_data["data"] and key["type"] == event_data["type"]: print("进入第一层判断") value = event_data['data'] for sub_key in key_path: print("进入循环") if sub_key in value.keys(): print("子键存在: ", sub_key) value = value[sub_key] else: value = None break if value is None: continue # 判断值类型,插入对应列 val_int = None val_str = None if type(value) == str: val_str = value elif type(value) == int: val_int = value # 添加event_data记录到会话 insert = Event_data(event_id=event["id"], type=event_data["type"], key=key["key_name"], value_int=val_int, value_str=val_str ) print("插入对象: ", insert) db.session.add(insert) def event_type_insert(event_type): event_type_list = getattr(g, "event_type_list", None) if event_type_list is None: check = [item.event_type for item in Event_type.query.all()] g.event_type_list = check event_type_list = check if event_type not in event_type_list: insert = Event_type(event_type=event_type) db.session.add(insert) g.event_type_list.append(event_type)
问题根源分析
- 连接池配置未生效:代码中配置的
pool_size=1和max_overflow=0与错误日志中的参数不符,说明应用实际使用的是SQLAlchemy默认连接池参数,大概率是配置加载顺序错误或环境变量覆盖了代码设置。 - 连接释放方式错误:手动调用
db.session.close()无法将连接正确归还到连接池,Flask-SQLAlchemy需要调用db.session.remove()或依赖框架自动回收连接。 - 低效查询占用连接:
insert_org等函数使用全表查询获取ID列表,数据量大时查询耗时久,且每个请求都会重复执行,高并发下快速耗尽连接池。 - 事务拆分不合理:同一个请求中多次调用
db.session.commit(),导致连接多次被占用,增加连接池压力。
解决方案
1. 修复连接池配置
调整初始化顺序确保配置生效,并根据并发量合理设置参数:
def create_app(): app = Flask(__name__, instance_relative_config=False) # 先加载配置对象,再设置URI避免覆盖 app.config.from_object('config.Config') db_uri = getenv("DB_URI") if db_uri: app.config['SQLALCHEMY_DATABASE_URI'] = db_uri with app.app_context(): db.init_app(app) from service.routes.post import post_event return app
修改连接池参数示例:
SQLALCHEMY_ENGINE_OPTIONS = { "pool_size": 10, "max_overflow": 20, "pool_recycle": 3600 # 避免连接超时失效 }
2. 正确释放连接
移除db.session.close(),改用db.session.remove():
# POST请求处理修改 try: # ... 插入逻辑 db.session.commit() insert_threads(org_id, site_id, event_dict, event_data) db.session.commit() db.session.remove() return {"result": "success"}, 200 except Exception as e: db.session.rollback() db.session.remove() return {"result": "unsuccessful"}, 409
3. 优化查询逻辑
- 替换全表查询为单条记录判断:
def insert_org(org_id): if not Organization.query.filter_by(id=org_id).first(): insert_org = Organization(id=org_id, name="temp_name") db.session.add(insert_org)
- 使用数据库唯一约束+
ON CONFLICT DO NOTHING避免先查后插:
def insert_org(org_id): db.session.execute( Organization.__table__.insert() .values(id=org_id, name="temp_name") .on_conflict_do_nothing(index_elements=['id']) )
- 只查询需要的字段减少数据传输:
# 原代码:check = [org.id for org in Organization.query.all()] # 修改为: check = [org.id for org in Organization.query.with_entities(Organization.id)]
4. 合并事务
将同一个请求中的多次提交合并为一次,减少连接占用:
# POST请求处理修改 try: insert_event = Event(**event_dict) db.session.add(insert_event) # 先执行关联表插入,再一次性提交 insert_threads(org_id, site_id, event_dict, event_data) db.session.commit() db.session.remove() return {"result": "success"}, 200 except Exception as e: db.session.rollback() db.session.remove() return {"result": "unsuccessful"}, 409
内容的提问来源于stack exchange,提问作者Jared Hembrow
相关产品推荐
相关产品推荐

