在 Ubuntu 环境下使用 Python 对接国产金仓数据库(KingbaseES),核心在于安装并配置专用的 ksycopg2 驱动。该驱动基于 DB API 2.0 规范,支持线程安全及异步通信。下面将分步骤演示从环境搭建到 CRUD 操作封装的全过程。
1. 环境准备与驱动安装
首先需要下载与当前 Python 版本匹配的 ksycopg2 驱动包。金仓官方提供了涵盖多个 Python 小版本的安装包,解压后通常包含对应版本的模块文件夹。
1.1 确认 Python 路径
上传驱动前,先确认 Python 的模块搜索路径:
import sys
print(sys.path)
运行后会输出类似 /usr/lib/python3/dist-packages 的路径,将解压后的 ksycopg2 文件夹复制至此目录即可。
1.2 配置环境变量
除了 Python 模块,还需将 KingbaseES 的 C 库文件路径添加到 LD_LIBRARY_PATH 中。若不确定库文件位置,可通过以下命令查找:
ps -ef | grep kingbase
假设找到路径为 /kingbase/data/KESRealPro/V009R002C012/Server/lib,则执行:
export LD_LIBRARY_PATH=/kingbase/data/KESRealPro/V009R002C012/Server/lib:$LD_LIBRARY_PATH
1.3 验证安装
尝试导入模块检查是否成功:
import ksycopg2
print("ksycopg2 驱动安装成功")
2. 连接数据库
建立连接需要提供数据库名、用户名、密码、主机地址和端口。建议将连接逻辑封装,便于后续复用。
import ksycopg2
def create_connection():
try:
conn = ksycopg2.connect(
database="TEST",
user="SYSTEM",
password="qwe123!@#",
host="127.0.0.1",
port="54321"
)
print("数据库连接成功")
return conn
except Exception as e:
print(f"连接数据库失败:{e}")
return None
connection = create_connection()
3. 创建数据表
在执行增删改查前,需确保测试表存在。这里以 user_info 表为例。
def create_table(conn):
try:
cursor = conn.cursor()
create_table_sql = """
CREATE TABLE IF NOT EXISTS user_info (
id INTEGER PRIMARY KEY,
username VARCHAR(50) NOT NULL,
age INTEGER,
created_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)
"""
cursor.execute(create_table_sql)
conn.commit()
cursor.close()
print("表创建成功")
except Exception as e:
print(f"创建表失败:{e}")
conn.rollback()
if connection:
create_table(connection)
4. 实现增删改查功能
4.1 新增数据
支持单条插入和批量插入。注意参数化查询可防止 SQL 注入。
def insert_data_example():
db_manager = KingbaseESManager(
dbname="test", user="SYSTEM", password="qwe123!@#",
host="127.0.0.1", port="54321"
)
if db_manager.connect():
# 插入单条数据
insert_sql = "INSERT INTO user_info (id, username, age) VALUES (%s, %s, %s)"
params = (1, '张三', 25)
db_manager.execute_update(insert_sql, params)
# 插入多条数据
users_to_insert = [
(2, '李四', 30), (3, '王五', 28),
(4, '赵六', 35), (5, '钱七', 22)
]
for user in users_to_insert:
db_manager.execute_update(insert_sql, user)
db_manager.disconnect()
insert_data_example()
4.2 查询数据
包括全量查询、条件查询及范围查询。
def query_data_example():
db_manager = KingbaseESManager(
dbname="test", user="SYSTEM", password="qwe123!@#",
host="127.0.0.1", port="54321"
)
if db_manager.connect():
# 查询所有用户信息
select_all_sql = "SELECT * FROM user_info ORDER BY id"
all_users = db_manager.execute_query(select_all_sql)
for user in all_users:
print(user)
# 根据 ID 查询特定用户
select_by_id_sql = "SELECT * FROM user_info WHERE id = %s"
user_by_id = db_manager.execute_query(select_by_id_sql, (2,))
print(user_by_id)
# 根据年龄范围查询
select_by_age_sql = "SELECT * FROM user_info WHERE age BETWEEN %s AND %s ORDER BY age"
users_by_age = db_manager.execute_query(select_by_age_sql, (25, 30))
for user in users_by_age:
print(user)
db_manager.disconnect()
query_data_example()
4.3 修改数据
更新操作同样需要提交事务。
def update_data_example():
db_manager = KingbaseESManager(
dbname="TEST", user="SYSTEM", password="your_password",
host="127.0.0.1", port="54321"
)
if db_manager.connect():
# 更新用户年龄
update_sql = "UPDATE user_info SET age = %s WHERE id = %s"
db_manager.execute_update(update_sql, (26, 1))
# 验证更新结果
select_sql = "SELECT * FROM user_info WHERE id = %s"
updated_user = db_manager.execute_query(select_sql, (1,))
print(updated_user)
db_manager.disconnect()
update_data_example()
4.4 删除数据
删除操作需注意事务回滚机制,确保数据安全。
def delete_data_example():
db_manager = KingbaseESManager(
dbname="TEST", user="SYSTEM", password="your_password",
host="127.0.0.1", port="54321"
)
if db_manager.connect():
# 删除 ID 为 5 的用户
delete_sql = "DELETE FROM user_info WHERE id = %s"
db_manager.execute_update(delete_sql, (5,))
# 验证删除结果
select_all_sql = "SELECT * FROM user_info ORDER BY id"
remaining_users = db_manager.execute_query(select_all_sql)
for user in remaining_users:
print(user)
db_manager.disconnect()
delete_data_example()
5. 封装管理类以便复用
为了减少重复代码,可以将上述逻辑封装为一个 KingbaseESManager 类。该类统一管理连接、游标及事务提交。
import ksycopg2
from datetime import datetime
class KingbaseESManager:
def __init__(self, dbname, user, password, host, port):
self.conn = None
self.db_params = {
"database": dbname,
"user": user,
"password": password,
"host": host,
"port": port
}
def connect(self):
"""连接数据库"""
try:
self.conn = ksycopg2.connect(**self.db_params)
print("数据库连接成功")
return True
except Exception as e:
print(f"连接数据库失败:{e}")
return False
def disconnect(self):
"""断开数据库连接"""
if self.conn:
self.conn.close()
print("数据库连接已关闭")
def execute_query(self, sql, params=None):
"""执行查询语句并返回结果"""
try:
cursor = self.conn.cursor()
if params:
cursor.execute(sql, params)
else:
cursor.execute(sql)
results = cursor.fetchall()
cursor.close()
return results
except Exception as e:
print(f"查询执行失败:{e}")
return None
def execute_update(self, sql, params=None):
"""执行更新操作(插入、更新、删除)"""
try:
cursor = self.conn.cursor()
if params:
cursor.execute(sql, params)
else:
cursor.execute(sql)
self.conn.commit()
affected_rows = cursor.rowcount
cursor.close()
print(f"操作成功,影响行数:{affected_rows}")
return affected_rows
except Exception as e:
self.conn.rollback()
print(f"操作执行失败:{e}")
return -1
def insert_user(self, id, username, age):
"""插入用户数据"""
sql = "INSERT INTO user_info (id, username, age) VALUES (%s, %s, %s)"
params = (id, username, age)
return self.execute_update(sql, params)
def select_all_users(self):
"""查询所有用户"""
sql = "SELECT * FROM user_info ORDER BY id"
return self.execute_query(sql)
def select_user_by_id(self, id):
"""根据 ID 查询用户"""
sql = "SELECT * FROM user_info WHERE id = %s"
params = (id,)
return self.execute_query(sql, params)
def update_user_age(self, id, new_age):
"""更新用户年龄"""
sql = "UPDATE user_info SET age = %s WHERE id = %s"
params = (new_age, id)
return self.execute_update(sql, params)
def delete_user(self, id):
"""删除用户"""
sql = "DELETE FROM user_info WHERE id = %s"
params = (id,)
return self.execute_update(sql, params)
def search_users_by_age_range(self, min_age, max_age):
"""根据年龄范围查询用户"""
sql = "SELECT * FROM user_info WHERE age BETWEEN %s AND %s ORDER BY age"
params = (min_age, max_age)
return self.execute_query(sql, params)
# 使用示例
if __name__ == "__main__":
db_manager = KingbaseESManager(
dbname="TEST", user="SYSTEM", password="qwe123!@#",
host="127.0.0.1", port="54321"
)
if db_manager.connect():
# 插入多条用户数据
users_to_insert = [
(1, '张三', 25), (2, '李四', 30),
(3, '王五', 28), (4, '赵六', 35),
(5, '钱七', 22)
]
for user in users_to_insert:
db_manager.insert_user(*user)
# 查询所有用户
print("所有用户信息:")
all_users = db_manager.select_all_users()
for user in all_users:
print(user)
# 根据 ID 查询用户
print("\n查询 ID 为 2 的用户:")
user_by_id = db_manager.select_user_by_id(2)
print(user_by_id)
# 更新用户年龄
print("\n更新张三的年龄为 26:")
db_manager.update_user_age(1, 26)
# 查询年龄在 25-30 岁之间的用户
print("\n年龄在 25-30 岁之间的用户:")
users_by_age = db_manager.search_users_by_age_range(25, 30)
for user in users_by_age:
print(user)
# 删除 ID 为 5 的用户
print("\n删除 ID 为 5 的用户:")
db_manager.delete_user(5)
# 再次查询所有用户
print("\n删除后的所有用户:")
remaining_users = db_manager.select_all_users()
for user in remaining_users:
print(user)
db_manager.disconnect()
6. 总结
在 Ubuntu 上通过 Python 操作 KingbaseES 主要涉及三个环节:驱动适配、环境配置与业务封装。安装时务必注意 Python 大版本的一致性,同时正确设置 LD_LIBRARY_PATH 以避免动态链接库加载错误。通过封装通用的管理类,可以有效降低维护成本,提升代码的可读性与复用性。


