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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 02:36:05