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

线程数超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()

关键细节说明

  1. Condition变量:threading.Condition结合了锁和等待队列的功能,能安全地让线程等待某个条件满足(这里就是有可用连接)。
  2. 循环等待:用while True而不是if,是为了避免虚假唤醒(线程可能在没有收到通知的情况下被唤醒),确保只有当连接真正可用时才继续执行。
  3. finally块:确保连接一定会被放回池,避免连接泄漏——如果线程在持有连接时抛出异常,没有finally的话连接就会丢失,导致池里的可用连接越来越少。

这种方式完全符合连接池的设计初衷:既控制了数据库连接的总数量(避免压垮数据库),又让线程优雅地排队等待,而不是直接崩溃。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 19:12:44