多线程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
相关产品推荐
相关产品推荐

