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

Python+psycopg2实现数据库重启后自动重连的最优方案

数据库单例类自动重连最优实现方案

我们在Python项目中使用如下数据库单例类,希望在数据库关闭并重启后自动重建连接,请问最优实现方式是什么?我们考虑在getconn函数中添加逻辑:先执行SELECT 1,捕获psycopg2异常后调用函数重置连接池。需注意我们的应用是多线程的,多个线程可同时访问该类,因此__new__函数中加入了锁机制。

当前使用的代码

import logging
import psycopg2
from psycopg2 import pool
from threading import Lock

# 假设全局锁定义在这里
dbConnLock = Lock()

class DBConnection:
    _instance = None
    logger: logging.Logger = None

    def __new__(cls, logger1: logging.Logger):
        dbConnLock.acquire()
        if cls._instance is None:
            cls._instance = object.__new__(cls)
            try:
                max_conn = 64
                keepalive_args = {"keepalives": 1, "keepalives_idle": 25, "keepalives_interval": 4,
                                  "keepalives_count": 9}
                print("starting creating pool")
                DBConnection._instance.pool = psycopg2.pool.ThreadedConnectionPool(5, 64, database='',
                                                                                   host='',
                                                                                   user='',
                                                                                   password='',
                                                                                   application_name='',
                                                                                   sslmode='require', connect_timeout=5,
                                                                                   **keepalive_args)
                logger1.info("Connection pool started")
                print("End creating pool")

            except Exception as ex:
                DBConnection._instance = None
                dbConnLock.release()
                raise ex
            cls._instance.__init__(logger1)
        dbConnLock.release()
        return cls._instance

    def __init__(self, logger1):
        self.logger = logger1

    def getconn(self, p_key):
        ps_connection = self._instance.pool.getconn(p_key)
        ps_cursor = ps_connection.cursor()
        ps_connection.autocommit = True
        self.logger.debug('Entering ' + str(p_key)+' '+ hex(id(ps_connection)))
        return ps_connection, ps_cursor

    def putconn(self, p_key, ps_connection, ps_cursor):
        if ps_cursor is not None:
            ps_cursor.close()
        if (ps_cursor is not None) and (ps_connection is not None):
            self._instance.pool.putconn(ps_connection, p_key)
        self.logger.debug('Exiting ' + str(p_key)+' '+hex(id(ps_connection)))

    def __del__(self):
        self._instance.pool.closeall()
        self.logger.info("Removing connection pool")

最优实现方案

核心思路

在连接获取阶段校验有效性,失效时触发线程安全的连接池重建,同时修正原有代码中的线程安全隐患与逻辑漏洞。

具体改进点

  1. 将全局锁改为类属性
    避免全局变量带来的耦合问题,把锁封装在类内部,更符合单例设计的封装性。

  2. 添加连接有效性校验
    在getconn中执行SELECT 1校验连接,捕获连接相关异常(如psycopg2.OperationalError),触发连接池重置。

  3. 实现线程安全的池重置方法
    单独编写_reset_pool方法,加锁确保同一时间只有一个线程重建池,避免多线程重复创建资源。

  4. 修正putconn逻辑漏洞
    原代码中判断ps_cursor is not None才放回连接,这会导致cursor为None时连接无法回收,改为优先判断connection是否有效。

  5. 细化异常处理
    只在连接失效类异常时触发重建,避免其他异常误触发池重置。

修改后的完整代码

import logging
import psycopg2
from psycopg2 import pool
from psycopg2 import OperationalError
from threading import Lock

class DBConnection:
    _instance = None
    logger: logging.Logger = None
    _lock = Lock()  # 类级锁,替代全局锁
    _pool_lock = Lock()  # 专门用于池重置的锁

    def __new__(cls, logger1: logging.Logger):
        with cls._lock:
            if cls._instance is None:
                cls._instance = object.__new__(cls)
                try:
                    cls._instance._init_pool(logger1)
                    cls._instance.__init__(logger1)
                except Exception as ex:
                    cls._instance = None
                    raise ex
        return cls._instance

    def __init__(self, logger1):
        if not hasattr(self, 'logger'):  # 避免重复初始化
            self.logger = logger1

    def _init_pool(self, logger1):
        """初始化连接池的私有方法"""
        max_conn = 64
        keepalive_args = {"keepalives": 1, "keepalives_idle": 25, "keepalives_interval": 4,
                          "keepalives_count": 9}
        logger1.info("Starting connection pool creation")
        self.pool = psycopg2.pool.ThreadedConnectionPool(
            minconn=5,
            maxconn=64,
            database='',
            host='',
            user='',
            password='',
            application_name='',
            sslmode='require',
            connect_timeout=5,
            **keepalive_args
        )
        logger1.info("Connection pool initialized successfully")

    def _reset_pool(self):
        """线程安全的连接池重置方法"""
        with self._pool_lock:
            # 先关闭旧池的所有连接
            if hasattr(self, 'pool'):
                try:
                    self.pool.closeall()
                    self.logger.info("Old connection pool closed")
                except Exception as ex:
                    self.logger.error(f"Failed to close old pool: {str(ex)}")
            # 重新初始化池
            self._init_pool(self.logger)

    def getconn(self, p_key):
        while True:
            try:
                ps_connection = self.pool.getconn(p_key)
                # 校验连接有效性
                with ps_connection.cursor() as cursor:
                    cursor.execute("SELECT 1")
                    cursor.fetchone()
                ps_connection.autocommit = True
                self.logger.debug(f'Entering {p_key} {hex(id(ps_connection))}')
                return ps_connection, ps_connection.cursor()
            except OperationalError as ex:
                self.logger.error(f"Connection invalid, resetting pool: {str(ex)}")
                # 重置连接池
                self._reset_pool()
                # 把无效连接放回池(后续池重置会关闭它)
                if 'ps_connection' in locals():
                    try:
                        self.pool.putconn(ps_connection, p_key)
                    except Exception as put_ex:
                        self.logger.error(f"Failed to put back invalid connection: {str(put_ex)}")
            except Exception as ex:
                self.logger.error(f"Unexpected error when getting connection: {str(ex)}")
                # 非连接异常直接抛出
                if 'ps_connection' in locals():
                    try:
                        self.pool.putconn(ps_connection, p_key)
                    except Exception as put_ex:
                        self.logger.error(f"Failed to put back connection: {str(put_ex)}")
                raise ex

    def putconn(self, p_key, ps_connection, ps_cursor):
        try:
            if ps_cursor is not None:
                ps_cursor.close()
            if ps_connection is not None:
                self.pool.putconn(ps_connection, p_key)
            self.logger.debug(f'Exiting {p_key} {hex(id(ps_connection))}')
        except Exception as ex:
            self.logger.error(f"Failed to put back connection: {str(ex)}")

    def __del__(self):
        if hasattr(self, 'pool'):
            try:
                self.pool.closeall()
                self.logger.info("Connection pool closed on instance deletion")
            except Exception as ex:
                self.logger.error(f"Failed to close pool on deletion: {str(ex)}")

关键细节说明

  • 双重锁机制:_lock用于单例实例创建,_pool_lock用于池重置,避免锁粒度太大影响性能。
  • 循环重试:getconn中使用循环,确保重置池后能重新获取有效连接。
  • 连接回收:即使连接无效,也尝试放回池,避免连接泄漏;重置池时会统一关闭所有旧连接。
  • 避免重复初始化:__init__中添加判断,防止单例实例重复初始化属性。

内容的提问来源于stack exchange,提问作者postgresuser

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 00:06:23