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

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

可能原因

  1. 连接资源耗尽:高频创建/关闭连接导致数据库连接数达上限,新连接无法建立,try块在psycopg2.connect阶段抛出异常,后续代码未执行。
  2. 事务阻塞:psycopg2默认开启事务,高频调用下未提交的事务堆积引发锁等待,导致查询无法执行。
  3. 异常日志不完整:当前仅打印异常信息,未记录异常发生的具体阶段,无法定位是连接、执行查询还是其他步骤出错。

修复方案

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 10:02:09