Gunicorn部署Flask应用遇MySQL连接超时及事务回滚问题求助
我有一个Python Flask应用,支持上传CSV文件并根据CSV文件表头创建GraphQL接口,通过向GraphQL端点发送POST请求及GraphQL查询来获取数据。应用使用Gunicorn运行(配置4个worker),需要从RDS MySQL实例(最大连接数146)获取文件表头信息,采用Flask-SQLAlchemy连接数据库。
部署到Kubernetes Pod后,首次发送POST请求到GraphQL端点一切正常,但一天后再次请求时出现两类错误:
First error:
Exception on /graphql/csv_source_service/v1/apps/8aa1eed7-9f86-4864-b163-a0b4c1eed8b3/csv_files/MjA2LXJlY29yZHMtZmlsZS0wMS5jc3Y= [POST]
Traceback (most recent call last):
File "/home/appuser/.local/lib/python3.9/site-packages/pymysql/connections.py", line 756, in _write_bytes
self._sock.sendall(data)
TimeoutError: [Errno 110] Connection timed out
second error:
Exception on /graphql/csv_source_service/v1/apps/8aa1eed7-9f86-4864-b163-a0b4c1eed8b3/csv_files/MjA2LXJlY29yZHMtZmlsZS0wMS5jc3Y= [POST]
Traceback (most recent call last):
File "/home/appuser/.local/lib/python3.9/site-packages/sqlalchemy/engine/base.py", line 1202, in _execute_context
conn = self._revalidate_connection()
File "/home/appuser/.local/lib/python3.9/site-packages/sqlalchemy/engine/base.py", line 469, in _revalidate_connection
raise exc.InvalidRequestError(
sqlalchemy.exc.InvalidRequestError: Can't reconnect until invalid transaction is rolled back
但使用Werkzeug开发服务器运行应用时不会出现此问题。
1. 修复全局Session/Engine的线程安全问题
Gunicorn多worker模式下,全局变量会被每个worker进程复制,当前代码的全局session管理逻辑会导致连接池混乱、跨请求复用无效连接。修改connection.py,移除全局变量,改用工厂函数管理连接和session:
from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker, scoped_session from sqlalchemy.pool import QueuePool import os host = '#####' username = '#####' password = '#####' database_name = '#####' engine_params = { 'poolclass': QueuePool, 'pool_size': 5, 'pool_pre_ping': True, 'echo': False, 'pool_recycle': 300, # 缩短连接回收时间,避免闲置连接被RDS/网络断开 'connect_args': { 'connect_timeout': 20, 'init_command': 'SET SESSION TRANSACTION ISOLATION LEVEL READ COMMITTED' } } def get_engine(): """创建并返回数据库引擎""" return create_engine( f'mysql+pymysql://{username}:{password}@{host}/{database_name}', **engine_params ) def get_session(): """为每个请求创建独立的scoped session""" engine = get_engine() session_factory = sessionmaker(autocommit=False, autoflush=False, bind=engine) return scoped_session(session_factory) def cleanup_session(session): """清理session,回滚未提交事务并释放连接""" if session: try: session.rollback() except Exception: pass session.remove()
2. 修正查询类的Session使用逻辑
当前CsvSourceSpAdapter在类初始化时就绑定全局session,会导致跨请求复用session引发事务问题。改为在方法内动态获取session,并确保请求结束后强制清理:
from connection import get_session, cleanup_session from sqlalchemy import and_ import logging logger = logging.getLogger(__name__) class CsvSourceSpAdapter: def get_by_id(self, app_id, file_id): session = get_session() try: result = session.query(CsvFile).filter( and_(CsvFile.file_id == file_id, CsvFile.app_id == app_id) ).first() if result: output_result = result.to_dict() logger.info("graphql stuff") logger.info(output_result) return { 'status': True, 'message': 'Found', 'data': output_result } else: logger.info("CSV file Not Found*******************************") return { 'status': False, 'message': 'CSV file not found' } finally: # 无论请求成功或失败,都必须清理session cleanup_session(session)
3. 添加Flask请求钩子自动管理Session
在Flask应用主文件中添加请求钩子,确保每个请求的session在请求结束后被正确清理,避免连接泄漏:
from flask import Flask, g from connection import get_session, cleanup_session app = Flask(__name__) @app.before_request def before_request(): """请求开始前绑定session到请求上下文""" g.db_session = get_session() @app.teardown_request def teardown_request(exception=None): """请求结束后清理session""" session = getattr(g, 'db_session', None) cleanup_session(session)
4. 关键参数说明
pool_recycle: 300:将连接回收时间从1小时改为5分钟,适配RDS或Kubernetes网络可能提前断开闲置连接的场景,让SQLAlchemy主动替换失效连接。pool_pre_ping: True:每次从连接池获取连接前,发送测试查询验证连接有效性,自动剔除失效连接。- 强制session清理:通过
finally块和请求钩子确保session被正确回滚、移除,避免无效事务残留。
内容的提问来源于stack exchange,提问作者Yohan Neranga

