Psycopg2中SELECT查询cursor.execute多次执行后出现异常问题
问题分析与修复方案
问题概述
以下Python函数每秒被调用一次,多次迭代后出现异常:无论传入什么参数,本该返回261行的SELECT查询,实际pd.DataFrame(cursor.fetchall())从未执行。问题并非每次请求触发,仅在高频调用多次后出现。
待排查代码
def GetTimeline(self, tenantId, tenantSiteId, query): connection = None timelineDataframe = pd.DataFrame() try: connection = psycopg2.connect(self.connectionString) print(f'[KeywordsTimelineRepository] connessione aperta') cursor = connection.cursor(cursor_factory = psycopg2.extras.RealDictCursor) querySQL="SELECT * FROM keywords_timeline WHERE tenantid = %s AND tenantsiteid = %s AND query = %s" data = (tenantId, tenantSiteId, query) cursor.execute(querySQL, data) timelineDataframe = pd.DataFrame(cursor.fetchall()) except (Exception, psycopg2.DatabaseError) as error: print(f'[KeywordsTimelineRepository]: {error}') finally: if connection is not None: connection.close() print(f'[KeywordsTimelineRepository] connessione chiusa') return timelineDataframe
可能原因
- 连接资源耗尽:高频创建/关闭连接导致数据库连接数达上限,新连接无法建立,
try块在psycopg2.connect阶段抛出异常,后续代码未执行。 - 事务阻塞:psycopg2默认开启事务,高频调用下未提交的事务堆积引发锁等待,导致查询无法执行。
- 异常日志不完整:当前仅打印异常信息,未记录异常发生的具体阶段,无法定位是连接、执行查询还是其他步骤出错。
修复方案
1. 补充详细日志定位问题
在关键步骤添加日志,包括参数、执行阶段、返回行数,同时打印异常栈信息:
import traceback def GetTimeline(self, tenantId, tenantSiteId, query): connection = None timelineDataframe = pd.DataFrame() try: print(f'[KeywordsTimelineRepository] Request params: tenantId={tenantId}, tenantSiteId={tenantSiteId}, query={query}') connection = psycopg2.connect(self.connectionString) print(f'[KeywordsTimelineRepository] connessione aperta') # 开启自动提交避免事务阻塞 connection.autocommit = True cursor = connection.cursor(cursor_factory = psycopg2.extras.RealDictCursor) querySQL="SELECT * FROM keywords_timeline WHERE tenantid = %s AND tenantsiteid = %s AND query = %s" data = (tenantId, tenantSiteId, query) print(f'[KeywordsTimelineRepository] Executing query with params: {data}') cursor.execute(querySQL, data) rows = cursor.fetchall() print(f'[KeywordsTimelineRepository] Fetched {len(rows)} rows') timelineDataframe = pd.DataFrame(rows) except (Exception, psycopg2.DatabaseError) as error: print(f'[KeywordsTimelineRepository] ERROR: {error}') traceback.print_exc() finally: if connection is not None: connection.close() print(f'[KeywordsTimelineRepository] connessione chiusa') print(f'[KeywordsTimelineRepository] Returning {len(timelineDataframe)} rows') return timelineDataframe
2. 使用连接池优化高频连接
高频调用下,单次创建/关闭连接效率极低,改用连接池复用连接:
import traceback import psycopg2.pool # 在类初始化方法中初始化连接池 def __init__(self): # 将connectionString解析为独立参数(示例,根据实际连接字符串调整) conn_params = { "dbname": "your_db", "user": "your_user", "password": "your_pwd", "host": "your_host", "port": "5432" } self.connection_pool = psycopg2.pool.SimpleConnectionPool( minconn=5, # 最小空闲连接数 maxconn=20, # 最大连接数 **conn_params ) def GetTimeline(self, tenantId, tenantSiteId, query): timelineDataframe = pd.DataFrame() connection = None try: connection = self.connection_pool.getconn() print(f'[KeywordsTimelineRepository] Got connection from pool') connection.autocommit = True cursor = connection.cursor(cursor_factory = psycopg2.extras.RealDictCursor) querySQL="SELECT * FROM keywords_timeline WHERE tenantid = %s AND tenantsiteid = %s AND query = %s" data = (tenantId, tenantSiteId, query) cursor.execute(querySQL, data) rows = cursor.fetchall() timelineDataframe = pd.DataFrame(rows) except (Exception, psycopg2.DatabaseError) as error: print(f'[KeywordsTimelineRepository] ERROR: {error}') traceback.print_exc() # 若连接出错,关闭并从池移除 if connection is not None: self.connection_pool.putconn(connection, close=True) connection = None finally: if connection is not None: self.connection_pool.putconn(connection) print(f'[KeywordsTimelineRepository] Returned connection to pool') return timelineDataframe
3. 检查数据库连接数配置
登录数据库执行以下命令查看连接数上限:
SHOW max_connections;
若当前调用频率(每秒1次)接近或超过上限,需调整数据库配置文件(如PostgreSQL的postgresql.conf)中的max_connections值,重启数据库生效。
内容的提问来源于stack exchange,提问作者P_R
相关产品推荐
相关产品推荐

