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

多线程Python应用PostgreSQL连接池耗尽问题求解

解决pool exhausted错误的方案

问题根源

你的连接池最大连接数设置为5,但update_records方法会为每条记录创建一个线程,线程总数远超过连接池容量。所有线程同时竞争有限的连接,当池内连接被全部占用后,新线程无法获取连接,就会抛出pool exhausted错误。此外,原代码还存在事务未提交、异常时连接未归还、无并发控制等问题,进一步加剧了这个问题。


具体解决步骤

1. 优化连接池配置与线程并发控制

  • 合理调整连接池大小:根据数据库允许的最大连接数(PostgreSQL默认100),适当调大maxconn,比如设为20,但不要超过数据库的max_connections参数。
  • 限制并发线程数:用threading.Semaphore控制同时运行的线程数,使其不超过连接池的最大连接数,避免无限制竞争连接。

2. 修复数据库操作的关键错误

  • 事务处理:INSERT/UPDATE操作后必须提交事务,否则修改不会生效;异常时要回滚事务,避免脏数据。
  • 安全归还连接:用try...finally块确保无论操作成功还是失败,连接都能归还到池里。
  • 避免空值报错:不要盲目调用fetchone()[0],UPDATE/INSERT操作通常返回受影响行数而非查询结果,需针对性处理。

3. 增强异常处理

获取连接失败时不要仅打印错误,应抛出异常让调用方处理,避免后续使用None连接导致崩溃。


修改后的代码示例

数据库连接池类(Database)

import psycopg2
from psycopg2 import pool

class Database(object):
    _instance = None

    def __new__(cls, *args, **kwargs):
        if cls._instance is None:
            cls._instance = super().__new__(cls)
            cls._instance.initialize_pool(*args, **kwargs)
        return cls._instance

    def initialize_pool(self, *args, **kwargs):
        # 调整maxconn为合理值,不超过数据库max_connections
        self.db_pool = pool.ThreadedConnectionPool(
            minconn=1,
            maxconn=20,
            *args, **kwargs
        )
        self.autocommit = True

    def get_connection(self):
        try:
            conn = self.db_pool.getconn()
            if self.autocommit:
                conn.autocommit = True
            return conn
        except Exception as e:
            raise RuntimeError(f"获取连接失败: {str(e)}") from e

    def return_connection(self, conn):
        try:
            self.db_pool.putconn(conn)
        except Exception as e:
            print(f"归还连接失败: {str(e)}")

    def insert_query(self, query, params=None):
        conn = None
        try:
            conn = self.get_connection()
            cursor = conn.cursor()
            cursor.execute(query, params or ())
            # 仅当使用RETURNING子句时才获取返回值
            result = cursor.fetchone()[0] if cursor.rowcount > 0 else None
            return result
        except Exception as e:
            if conn:
                conn.rollback()
            raise e
        finally:
            if conn:
                self.return_connection(conn)

    def update_query(self, query, params=None):
        conn = None
        try:
            conn = self.get_connection()
            cursor = conn.cursor()
            cursor.execute(query, params or ())
            return cursor.rowcount  # 返回受影响行数,更适合UPDATE操作
        except Exception as e:
            if conn:
                conn.rollback()
            raise e
        finally:
            if conn:
                self.return_connection(conn)

主应用类(MainApp)

from .db import Database
import threading

class MainApp:
    def __init__(self):
        self.records = [1, 2, 3, ...]  # 你的记录ID列表
        # 用参数而非连接字符串更安全清晰
        self.conn_params = {
            "dbname": "你的数据库名",
            "user": "用户名",
            "password": "密码",
            "host": "数据库地址"
        }
        self.db = Database(**self.conn_params)
        # 并发数等于连接池最大连接数,避免连接耗尽
        self.semaphore = threading.Semaphore(self.db.db_pool.maxconn)

    def insert_into_db(self, query, params=None):
        def worker():
            with self.semaphore:
                try:
                    self.db.insert_query(query, params)
                except Exception as e:
                    print(f"插入失败: {str(e)}")
        thread = threading.Thread(target=worker)
        thread.start()

    def update_records(self, query_template):
        def worker(record_id):
            with self.semaphore:
                # 使用参数化查询避免SQL注入,这是安全最佳实践
                try:
                    affected_rows = self.db.update_query(
                        query_template,
                        (record_id,)
                    )
                    print(f"更新记录{record_id}完成,影响行数: {affected_rows}")
                except Exception as e:
                    print(f"更新记录{record_id}失败: {str(e)}")
        
        for record in self.records:
            thread = threading.Thread(target=worker, args=(record,))
            thread.start()

# 使用示例
app = MainApp()
insert_query = "INSERT INTO your_table (col1) VALUES (%s) RETURNING id;"
app.insert_into_db(insert_query, ("测试值",))
update_query_template = "UPDATE your_table SET col2 = '已更新' WHERE id = %s;"
app.update_records(update_query_template)

额外注意事项

  • 参数化查询:永远不要用字符串拼接生成SQL,必须使用psycopg2的参数化方式,防止SQL注入。
  • 数据库连接限制:定期检查数据库的max_connections配置,确保连接池的maxconn不超过该值,预留部分连接给其他应用。
  • 线程池替代方案:如果记录量极大,建议使用concurrent.futures.ThreadPoolExecutor替代手动创建线程,更高效且易于管理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 04:55:10