线程数超psycopg2连接池上限时崩溃,如何实现等待获取连接?
解决psycopg2 ThreadedConnectionPool连接耗尽时的等待问题
你碰到的这个问题其实很常见——默认的ThreadedConnectionPool在连接耗尽时会直接抛出PoolError,而不会让线程等待可用连接。这并不是你的操作有误,只是这个池本身没有内置等待逻辑而已,咱们可以给它加一层封装来实现你想要的行为。
修复思路
我们可以继承原有的ThreadedConnectionPool,利用Python的threading.Condition来实现等待-通知机制:
- 当线程请求连接但池已耗尽时,让线程进入等待状态
- 当有线程把连接放回池时,唤醒一个等待的线程来获取连接
修改后的可复现代码
import threading import psycopg2 from psycopg2 import pool class WaitingThreadedConnectionPool(pool.ThreadedConnectionPool): def __init__(self, minconn, maxconn, *args, **kwargs): super().__init__(minconn, maxconn, *args, **kwargs) # 初始化条件变量,用于线程间的等待和通知 self._condition = threading.Condition() def getconn(self, key=None): with self._condition: # 循环尝试获取连接,直到成功 while True: try: # 调用父类的getconn方法 return super().getconn(key) except psycopg2.pool.PoolError: # 连接池耗尽,当前线程进入等待 self._condition.wait() def putconn(self, conn, key=None, close=False): with self._condition: # 调用父类的putconn方法放回连接 super().putconn(conn, key, close) # 通知一个等待的线程:现在有可用连接了 self._condition.notify() # 初始化带等待机制的连接池,最小1个连接,最大10个 conn_pool = WaitingThreadedConnectionPool( 1, 10, host='127.0.0.1', user='john', password='1234', dbname='test', port=1234 ) class Foo(threading.Thread): def __init__(self): super().__init__() def run(self): global conn_pool conn = None try: # 现在调用getconn会自动等待,直到有可用连接 conn = conn_pool.getconn() cur = conn.cursor() sql_query = "SELECT id from test_table;" cur.execute(sql_query) # 原代码print(cur.execute(...))会输出None,这里改成获取实际结果 result = cur.fetchone() print(f"线程 {self.name} 查询结果: {result}") cur.close() finally: # 用finally确保无论是否出错,连接都会被放回池 if conn: conn_pool.putconn(conn) num_threads = 20 threads = [] for i in range(num_threads): threads.append(Foo()) for thread in threads: thread.start() for thread in threads: thread.join() conn_pool.closeall()
关键细节说明
- Condition变量:
threading.Condition结合了锁和等待队列的功能,能安全地让线程等待某个条件满足(这里就是有可用连接)。 - 循环等待:用
while True而不是if,是为了避免虚假唤醒(线程可能在没有收到通知的情况下被唤醒),确保只有当连接真正可用时才继续执行。 - finally块:确保连接一定会被放回池,避免连接泄漏——如果线程在持有连接时抛出异常,没有finally的话连接就会丢失,导致池里的可用连接越来越少。
这种方式完全符合连接池的设计初衷:既控制了数据库连接的总数量(避免压垮数据库),又让线程优雅地排队等待,而不是直接崩溃。
内容的提问来源于stack exchange,提问作者puk
相关产品推荐
相关产品推荐

