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")
最优实现方案
核心思路
在连接获取阶段校验有效性,失效时触发线程安全的连接池重建,同时修正原有代码中的线程安全隐患与逻辑漏洞。
具体改进点
将全局锁改为类属性
避免全局变量带来的耦合问题,把锁封装在类内部,更符合单例设计的封装性。添加连接有效性校验
在getconn中执行SELECT 1校验连接,捕获连接相关异常(如psycopg2.OperationalError),触发连接池重置。实现线程安全的池重置方法
单独编写_reset_pool方法,加锁确保同一时间只有一个线程重建池,避免多线程重复创建资源。修正putconn逻辑漏洞
原代码中判断ps_cursor is not None才放回连接,这会导致cursor为None时连接无法回收,改为优先判断connection是否有效。细化异常处理
只在连接失效类异常时触发重建,避免其他异常误触发池重置。
修改后的完整代码
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
相关产品推荐
相关产品推荐

