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

使用Python函数连接并操作PostgreSQL的dvdrental数据库

改造PostgreSQL连接函数并实现查询结果转Pandas DataFrame

需求概述

作为Python与PostgreSQL新手,已完成以下基础工作:

  • 创建.config文件夹及存储登录凭证的database.ini文件
  • 在src目录编写config.py,通过ConfigParser读取配置
  • 在src目录编写tasks.py,实现基础connect函数可打印数据库版本

需要完成:

  • 修改connect函数使其通用化,返回数据库连接与游标,支持后续查询操作
  • 实现查询函数(如select_from_table),并将查询结果转换为Pandas DataFrame

修改后的connect函数

原connect函数会自动关闭连接和游标,无法复用连接进行后续查询。修改后保留连接与游标,交由调用者管理生命周期:

import pandas as pd
from clients.config import config
import psycopg

def connect():
    """连接PostgreSQL数据库,返回连接对象与游标"""
    conn = None
    cur = None
    try:
        # 读取连接参数
        params = config()
        print('连接PostgreSQL数据库中...')
        conn = psycopg.connect(**params)
        # 创建游标
        cur = conn.cursor()
        # 验证连接成功
        cur.execute('SELECT version()')
        db_version = cur.fetchone()
        print(f"PostgreSQL版本: {db_version[0]}")
        return conn, cur
    except (Exception, psycopg.DatabaseError) as error:
        print(f"连接出错: {error}")
        # 出错时若已创建连接则关闭
        if conn is not None:
            conn.close()
        return None, None

实现查询函数并转换为DataFrame

以下查询函数支持指定schema和表名,将查询结果转为带列名的Pandas DataFrame:

def select_from_table(cursor, table_name, schema):
    """查询指定schema下的表,返回结果列表与列名"""
    try:
        # 设置搜索路径
        cursor.execute(f"SET search_path TO {schema}, public;")
        # 执行查询
        cursor.execute(f"SELECT * FROM {table_name};")
        # 获取查询结果与列名
        results = cursor.fetchall()
        columns = [desc[0] for desc in cursor.description]
        return results, columns
    except (Exception, psycopg.DatabaseError) as error:
        print(f"查询出错: {error}")
        return None, None

def query_to_dataframe(table_name, schema):
    """封装查询流程,直接返回Pandas DataFrame"""
    conn, cur = connect()
    if cur is None:
        return pd.DataFrame()
    
    results, columns = select_from_table(cur, table_name, schema)
    # 转换为DataFrame
    df = pd.DataFrame(results, columns=columns)
    
    # 关闭游标与连接
    cur.close()
    conn.close()
    return df

使用示例

if __name__ == '__main__':
    # 方式一:手动管理连接与游标
    conn, cur = connect()
    if conn and cur:
        # 执行查询
        results, columns = select_from_table(cur, "your_table", "your_schema")
        # 转换为DataFrame
        df = pd.DataFrame(results, columns=columns)
        print(df.head())
        
        # 用完务必关闭资源
        cur.close()
        conn.close()
        print("数据库连接已关闭")
    
    # 方式二:使用封装好的函数直接获取DataFrame
    # df = query_to_dataframe("your_table", "your_schema")
    # print(df.head())

注意事项

  1. SQL注入风险:上述示例使用字符串格式化拼接SQL语句,仅适用于测试环境。生产环境请对表名/ schema名做白名单验证,避免注入风险。
  2. 连接管理:使用完连接与游标后务必关闭,避免资源泄漏;也可使用上下文管理器自动管理生命周期:
    with psycopg.connect(**config()) as conn:
        with conn.cursor() as cur:
            cur.execute("SELECT * FROM your_table;")
            results = cur.fetchall()
    
  3. 异常处理:实际使用中可根据需求扩展异常处理逻辑,比如增加连接重试机制。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 07:10:59