使用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())
注意事项
- SQL注入风险:上述示例使用字符串格式化拼接SQL语句,仅适用于测试环境。生产环境请对表名/ schema名做白名单验证,避免注入风险。
- 连接管理:使用完连接与游标后务必关闭,避免资源泄漏;也可使用上下文管理器自动管理生命周期:
with psycopg.connect(**config()) as conn: with conn.cursor() as cur: cur.execute("SELECT * FROM your_table;") results = cur.fetchall() - 异常处理:实际使用中可根据需求扩展异常处理逻辑,比如增加连接重试机制。
内容的提问来源于stack exchange,提问作者jenna
相关产品推荐
相关产品推荐

