资讯详情

资讯详情

Python+FastAPI连接KingbaseES:从裸SQL到连接池与事务管理的完整实践

我一直跟团队强调一句话用 Python 给 KingbaseES 写服务别再用裸 SQL 直连了。上个月客户现场的报表接口把生产库连接数冲到 300DBA 半夜打电话过来我打开代码一看SQL 全是用字符串拼接在 FastAPI 视图函数里没有连接池没有事务控制连查询参数都是直接塞进 SQL 的。这就是典型的“裸 SQL 连金仓”的连环雷。这篇文章是我用 Python FastAPI 给 KingbaseES 搭一套可上线 API 的完整复盘从驱动选型到连接池、事务、慢 SQL 排查都覆盖了。适合正在做数据库迁移、或者刚拿到金仓环境准备写后端服务的 Python 工程师看完可以直接照着搭。1. 裸 SQL 连金仓的四个典型事故1.1 SQL 注入不是拼一个参数那么简单很多人写金仓的查询接口时会下意识写出这种代码cursor.execute(fSELECT * FROM user_info WHERE login_name {login_name})单看一个接口好像没什么但login_name一旦由用户传进来就是一个天然的注入点。比如用户输入admin OR 11拼出来就成了SELECT * FROM user_info WHERE login_name admin OR 11表里所有账号都能被查出来。再狠一点的用; DROP TABLE ...这种语法直接改表结构虽然金仓的账号权限配置可能挡住了部分操作但一旦应用账号权限过大后果就是整个业务表没了。正确的姿势是用驱动提供的参数绑定接口cursor.execute( SELECT id, login_name, display_name FROM user_info WHERE login_name %s, (login_name,), )ksycopg2和psycopg2底层走的是预编译语句参数值会通过独立协议通道传给数据库跟 SQL 文本本身完全分离。这不是锦上添花而是安全底线。你可以在数据访问层自己定义一套统一的查询入口让所有 SQL 都必须走参数绑定而不是靠程序员自觉。1.2 连接不池化压测一到连接数先崩裸 SQL 连金仓最常见的问题还不是注入而是每次请求都新建一个物理连接。TCP 握手、认证握手、后端分配进程和内存一套流程走下来比查询本身还慢。更麻烦的是数据库的max_connections是有限制的默认可能只有 100 或者 200报表系统一旦有人多点了几次刷新连接数瞬间打满其他服务跟着遭殃。你会在数据库日志里看到类似这样的错误FATAL: remaining connection slots are reserved for non-replication superuser connections第一反应是“数据库挂了”其实数据库活着是连接被占满了。这类事故在高并发压测中几乎 100% 会出现。解决办法就一个连接池。让一组连接反复复用而不是每次请求都重新创建这能在不提高数据库压力的情况下扛住更高的并发。1.3 事务没人管半截数据写进生产库裸 SQL 直连时代事务经常是乱写的conn get_connection() try: conn.cursor().execute(INSERT INTO order_header ...) conn.cursor().execute(INSERT INTO order_line ...) conn.commit() except Exception: conn.rollback()看起来有commit和rollback但实际业务里常常是 A 服务写头部、B 服务写明细甚至两个操作各拿各的连接。于是就会出现“订单头写了明细没写”这种半截数据。等到对账发现对不上已经在生产环境里躺了很久了。事务边界必须按照“一个业务操作 一个连接 一个事务”的粒度来设计。在我的实践里这个边界放在 Service 层而不是 Repository 层。后面第四章会详细展开。1.4 SQL 散落全项目表结构一改全局加班裸 SQL 连金仓的项目还有个典型特征SQL 字符串散落在路由、中间件、工具函数里。表结构发生变化比如user_info表把mobile字段改名为phone你要全局搜mobile然后一个文件一个文件改漏改任何一个线上接口就直接报column does not exist。这不是代码风格问题而是维护模型的问题。SQL 是跟表结构强耦合的资产它应该有明确的归属地。我用 Repository 层统一收编所有的 SQL路由层只负责参数校验和 HTTP 响应数据访问层只负责 SQL。这样表结构变更时需要动的地方很集中排查范围也小得多。2. 环境准备驱动选型和第一行能跑的连接代码2.1 Python 版本与虚拟环境做金仓的 API 服务Python 版本建议直接用 3.10 以上版本。如果现场机器上同时有多个 Python建议用虚拟环境隔离避免依赖互相污染。python3 -m venv .venv source .venv/bin/activate依赖方面基础安装这几样就够了pip install fastapi uvicorn[standard] pydantic-settings数据库驱动单独说下面是重点。2.2 KingbaseES 的 Python 驱动怎么选金仓的官方 Python 驱动常见的是ksycopg2API 基本复刻 psycopg2。也有一部分现场因为数据库配置了 PG 兼容线协议直接用psycopg2-binary也能连上。我的建议是根据环境区分方案优点缺点建议ksycopg2官方驱动跟金仓内核版本匹配出问题有据可查安装包通常要从安装介质获取不能完全依赖 pip 拉取生产环境首选psycopg2-binarypip 一行安装开发环境方便不是金仓官方测试覆盖范围部分小版本可能遇到兼容性坑开发环境应急可以上线前必须回归SQLAlchemy自带连接池、事务管理、SQL 编译能力金仓大体兼容 PG 方言但个别特性要现场验证团队愿意投入可以选否则不强行上这里要特别提醒一句ksycopg2不是随随便便pip install ksycopg2就能拿到正版官方包的网上存在同名仿包风险。最好从金仓官网或安装介质里拿对应 Python 版本的安装文件确认哈希后再安装。2.3 连接串和连通性自测金仓的默认端口常见是 54321不是 PostgreSQL 的 5432。我见过太多人拿着 PG 的习惯去连金仓连不上还以为是密码错了。先确认现场 DBA 给的端口。连接串可以这样组织import ksycopg2 conn ksycopg2.connect( host127.0.0.1, port54321, dbnameapp_db, userapp_rw, password******, connect_timeout5, ) cur conn.cursor() cur.execute(select version()) print(cur.fetchone()) cur.close() conn.close()能打印出版本信息说明这第一关过了。这里我把connect_timeout设成 5 秒是为了避免数据库不可达时 API 进程长时间卡在网络超时上。生产环境里这个参数必须显式设置否则一个挖断光缆的故障就能让所有请求挂在连接阶段。3. 项目骨架把 SQL 从路由层收编到 Repository3.1 一个能上线的目录结构我搭金仓 API 服务时用的目录结构大概是这样的app/ main.py config.py db.py routers/ user_router.py services/ user_service.py repositories/ user_repo.py schemas/ user_schema.py每一层职责很明确路由层只做参数校验和响应封装Service 层做业务规则、事务编排Repository 层只写 SQLdb.py统一管连接池。这个结构看起来朴素但很稳团队新成员过来也能一眼看懂。3.2 配置管理密码绝不能写死连接信息放在源码里等于裸奔。我用pydantic-settings从环境变量读取配置本地开发时用.env文件线上则通过部署平台注入环境变量。# config.py from pydantic_settings import BaseSettings class Settings(BaseSettings): kingbase_host: str 127.0.0.1 kingbase_port: int 54321 kingbase_user: str kingbase_password: str kingbase_db: str db_pool_min: int 2 db_pool_max: int 10 property def dsn(self) - str: return ( fhost{self.kingbase_host} port{self.kingbase_port} fdbname{self.kingbase_db} user{self.kingbase_user} fpassword{self.kingbase_password} connect_timeout5 ) settings Settings()注意.env文件一定要写进.gitignore不然一个团队多人协作密码迟早会被推到 Git 仓库里。如果公司有 Vault 这类密钥管理平台优先用平台注入而不是本地文件。3.3 Repository 模式到底解决什么问题Repository 模式不是高大上的架构概念它解决的就是 SQL 散落各处的实际问题。Repository 层把 SQL 从业务代码里剥离出来路由层和 Service 层都不再出现 SQL 片段。比如一个用户查询Repository 里是这样# repositories/user_repo.py class UserRepo: staticmethod def get_by_login_name(conn, login_name: str): sql SELECT id, login_name, display_name FROM user_info WHERE login_name %s with conn.cursor() as cur: cur.execute(sql, (login_name,)) return cur.fetchone()这个conn是 Service 层传进来的不是 Repository 自己创建的。这么做有个关键好处多个 Repository 方法可以共享同一个事务连接而不会出现每个方法各拿一个连接导致事务失效的问题。4. 核心实现一个 DatabasePool 管住连接、事务和参数绑定4.1 DatabasePool 实现连接池我用ksycopg2.pool.ThreadedConnectionPool实现一个简单的封装# db.py import ksycopg2 from ksycopg2 import pool from contextlib import contextmanager class DatabasePool: def __init__(self, dsn: str, minconn: int 2, maxconn: int 10): self._pool pool.ThreadedConnectionPool(minconn, maxconn, dsn) contextmanager def connection(self): conn self._pool.getconn() try: yield conn conn.commit() except Exception: conn.rollback() raise finally: self._pool.putconn(conn) def close(self): self._pool.closeall()这里需要理解一个关键点ThreadedConnectionPool是线程级别的连接池。FastAPI 的同步接口默认跑在线程池里所以这个池子刚好匹配。如果你的接口全部写成async def又在里面直接执行阻塞式 SQL那就不是连接池能解决的问题而是把整个事件循环卡死了。所以我的原则是数据库操作是同步的接口就老老实实写成同步def让 FastAPI 自己管理线程池。千万别为了“显得更快”写成async def结果把性能拖垮。4.2 生命周期FastAPI 启动时建池关闭时销毁连接池要跟服务生命周期绑定。FastAPI 现在推荐用 lifespan 管理# main.py from contextlib import asynccontextmanager from fastapi import FastAPI, Request asynccontextmanager async def lifespan(app: FastAPI): app.state.db_pool DatabasePool(settings.dsn, settings.db_pool_min, settings.db_pool_max) yield app.state.db_pool.close() app FastAPI(lifespanlifespan)依赖注入时从app.state里取连接池def get_db_pool(request: Request) - DatabasePool: return request.app.state.db_pool启动时只初始化一次关闭时统一销毁这是最不容易泄漏资源的写法。4.3 事务边界放在 Service 层而不是 Repository前面说过Repository 方法不应该自己拿连接而是接收外部传入的conn。事务边界放在 Service 层用with db_pool.connection()包住整个业务操作# services/user_service.py class UserService: def __init__(self, db_pool: DatabasePool): self.db_pool db_pool def get_or_create_user(self, login_name: str, display_name: str): with self.db_pool.connection() as conn: user UserRepo.get_by_login_name(conn, login_name) if user: return user return UserRepo.create_user(conn, login_name, display_name)在这个上下文里无论中间哪个 SQL 抛异常connection()都会执行rollback()所有写操作全部回滚。只有全部成功才会走到commit()。这里容易踩的坑是有人觉得 Repository 方法内部自己connect()、自己commit()也行。结果两个写操作一个成功一个失败数据就裂了。所以事务边界必须往上提提到业务语义完整的 Service 层。4.4 参数绑定SQL 注入最直接的那道闸在金仓的 Python 驱动里参数占位符是%s不是?也不是 SQLAlchemy 的:name。规则很简单值必须通过参数传不允许拼字符串只有一个参数也要传元组比如(login_name,)像%这种模运算符号出现在 SQL 里要写成%%转义看几个常见的写法# 等值查询 cur.execute( SELECT * FROM user_info WHERE login_name %s, (login_name,), ) # LIKE 查询通配符放在值里不放在 SQL 模板里 cur.execute( SELECT * FROM user_info WHERE display_name LIKE %s, (f%{keyword}%,), ) # IN 查询用 ANY 数组参数避免动态拼占位符 cur.execute( SELECT * FROM order_info WHERE status ANY(%s), (status_list,), )注意一点表名、列名这类标识符不能参数化绑定。如果业务上必须动态选择列名或表名一定要用白名单校验不要直接信任前端传的值。ALLOWED_COLUMNS {login_name, display_name, created_at} if order_by not in ALLOWED_COLUMNS: raise HTTPException(status_code400, detailinvalid order_by)参数绑定是一道闸白名单是第二道闸两道闸都落下SQL 注入基本没戏。4.5 同步驱动放进异步框架的正确姿势金仓官方驱动目前主要还是同步 DBAPI 风格。FastAPI 里用同步驱动有两条路一条是接口写成同步defFastAPI 自动放到线程池执行app.get(/users/{login_name}) def get_user(login_name: str, db_pool: DatabasePool Depends(get_db_pool)): service UserService(db_pool) return service.get_by_login_name(login_name)另一条是接口必须写成async def用run_in_threadpool手动丢到线程池from starlette.concurrency import run_in_threadpool app.get(/users/{login_name}) async def get_user(login_name: str, db_pool: DatabasePool Depends(get_db_pool)): service UserService(db_pool) return await run_in_threadpool(service.get_by_login_name, login_name)两条路能跑通但我实际项目里默认走第一条。代码更短也不容易犯“在事件循环里休眠线程”的低级错误。5. 上线前必须过的几道检查关5.1 连接数预算要按 worker 数乘一遍连接池不是越大越好。很多人会犯一个错单机连接池设了 50看起来不多结果服务用 Gunicorn 起了 4 个 worker每个 worker 50 个连接瞬间就是 200 个连接压在数据库上。一个可用的估算公式预估连接数 worker 进程数 × 单进程连接池上限 数据库剩余连接数 max_connections × 0.7 - 运维预留连接数比如金仓max_connections 200运维预留 40那可用连接大概是 1004 个 worker 的话每个 worker 的连接池上限控制在 25 以内会比较安全。此外连接池里的连接空闲久了可能被数据库回收。解决这个问题最土也最有效的办法是健康检查时执行SELECT 1强制确认连接是活的app.get(/healthz) def healthz(request: Request): db_pool request.app.state.db_pool with db_pool.connection() as conn: conn.cursor().execute(SELECT 1) return {status: ok}5.2 慢 SQL 先抓到再谈优化上线前一定要把慢 SQL 日志打开。金仓的 PostgreSQL 兼容性做得比较好很多参数跟 PG 类似。可以找 DBA 帮忙设置log_min_duration_statement 1000超过 1 秒的 SQL 都写到日志里。排查时我会先用EXPLAIN ANALYZE看执行计划EXPLAIN (ANALYZE, BUFFERS) SELECT id, login_name, display_name FROM user_info WHERE login_name test;看三点是不是走了索引、有没有全表扫描、扫描行数和返回行数是否差异巨大。如果Seq Scan出现在高频查询上且表行数已经上万就果断建索引。另外金仓的统计视图名在不同版本里可能不一样我在现场见过pg_stat_activity也见过sys_stat_activity。写监控脚本前先确认一下当前实例到底支持哪个不要想当然。5.3 异常别裸奔错误码映射数据库异常直接抛给前端除了暴露表结构信息体验也差。我通常会在全局异常处理器里做一层映射数据库异常HTTP 状态码业务含义UniqueViolation409唯一键冲突数据已存在ForeignKeyViolation400关联记录不存在ConnectionException503数据库连接不可用QueryCanceled/ 超时504SQL 执行超时核心思路是驱动抛出的具体异常在 Service 层被捕获转换成展现层能理解的状态码和消息。不要让原始异常堆栈直接打到响应体里。from fastapi.responses import JSONResponse from ksycopg2.errors import UniqueViolation app.exception_handler(UniqueViolation) async def unique_violation_handler(request, exc): return JSONResponse(status_code409, content{detail: 数据已存在})5.4 最小权限账号应用账号别用管理员给 API 用的数据库账号应该只授权业务必须的 DML 操作不要给 DDL 权限更不要直接拿高权限账号跑应用。CREATE USER app_rw WITH PASSWORD ******; GRANT CONNECT ON DATABASE app_db TO app_rw; GRANT USAGE ON SCHEMA public TO app_rw; GRANT SELECT, INSERT, UPDATE, DELETE ON ALL TABLES IN SCHEMA public TO app_rw;如果只是报表查询可以再单独建一个只读账号。权限控制这种事平时没什么存在感真出事的时候就是最后一道防线。6. 这些坑是我在生产环境里踩出来的6.1 多进程 worker 让连接数直接翻倍有一次我在客户现场部署本机压测一切正常一上生产环境就报连接数不足。排查下来发现是 Gunicorn 起了 4 个 worker每个 worker 里 FastAPI 应用都独立创建了一个连接池实际连接数是单机压测的 4 倍。这个问题不是代码 bug而是部署架构和连接池配置没联动起来。后来我在部署文档里写死了计算公式worker 数 × pool_max ≤ max_connections × 0.5。6.2 中文乱码十次里有八次是 client_encoding有段时间接口返回的中文偶尔变乱码查了很久最后发现是客户端连接的client_encoding没有显式指定。数据库本身编码是 UTF8但连接参数没带client_encodingutf8换了一个环境之后编码就飘了。解决方式是在 DSN 里直接加上property def dsn(self) - str: return ( fhost{self.kingbase_host} port{self.kingbase_port} fdbname{self.kingbase_db} user{self.kingbase_user} fpassword{self.kingbase_password} connect_timeout5 fclient_encodingutf8 )6.3 LIMIT 参数化的类型坑参数绑定不是万灵药。有一次写分页接口把页码和页大小都传进参数结果报错说参数类型不匹配。原因是LIMIT %s里的参数必须整数前端传的字符串明明能转成 10驱动却偏不帮你转。# 这样会报错 page_size 10 cur.execute(SELECT * FROM order_info LIMIT %s, (page_size,)) # 得转成 int cur.execute(SELECT * FROM order_info LIMIT %s, (int(page_size),))统一收口在 Repository 层之后这类问题只在少数几个函数里出现排查范围小很多。6.4 官方驱动版本要跟数据库小版本对表金仓有 V8、V9 等不同大版本同一大版本下面还有小版本差异。官方驱动的发布节奏不一定能跟上所有数据库小版本所以现场如果报驱动兼容性错误第一件事不是去改代码而是去问 DBA 当前数据库内核小版本是多少再去官网找匹配的驱动版本。我之前就遇到过一个ksycopg2连不上特定小版本实例的情况换回同批次发布的驱动就好了。6.5 健康检查也可能变成拖累健康检查本意是好的但如果不加约束比如每台机器每秒钟都打一次/healthz每次健康检查都去连接池拿一个连接执行SELECT 1几百个服务实例叠加起来也是一笔不小的查询量。更关键的是如果连接池里的连接已经失效健康检查会反复触发重连。我的做法是健康检查频率控制在 10 秒以上且只做最简单的SELECT 1不要在里面查业务大表。如果让我现在从头搭一遍金仓 API我会把上面这些检查项写成 CI 里的一个脚本拉代码后先跑单元测试再自动跑一遍连接池压测和慢 SQL 检查应用账号只保留业务需要的权限高权限账号只留在运维手里。跑通第一个接口很容易难的是保证服务上线三个月后数据库朋友的电话不要在半夜响起来。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →