首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >别再裸 SQL 连金仓了:用 Python + FastAPI 给 KingbaseES 搭一套可上线的 API

别再裸 SQL 连金仓了:用 Python + FastAPI 给 KingbaseES 搭一套可上线的 API

原创
作者头像
小白学大数据
发布2026-09-04 17:10:39
发布2026-09-04 17:10:39
110
举报

前言:为什么用 Python 给金仓"搭接口"金仓数据库(KingbaseES,下文简称 KES)是电科金仓推出的企业级关系型数据库,在政务、金融、能源、交通等关键行业落地极广,也是当前"信创"(信息技术应用创新)国产化替代浪潮中的主流选型之一。它最讨开发者喜欢的一点是 对 PostgreSQL 通信协议的高度兼容——这意味着 PostgreSQL 生态里那套成熟的 Python 工具链(驱动、ORM、连接池)几乎可以"免适配"地平移过来,迁移成本极低。本文的目标很明确:用一套工程化、可上线的方式,基于 Python 把 KES 封装成一个标准的 RESTful API。与网上常见的"Flask + 裸 SQL"示例不同,我们采用更贴合生产实践的技术组合:FastAPI:异步优先、自带 OpenAPI 文档、依赖注入与类型校验开箱即用;SQLAlchemy 2.0 ORM:以 Python 类声明表结构,从根上消灭字符串拼接 SQL;ksycopg2:金仓官方提供的 Python 驱动,API 与 psycopg2 高度一致;Pydantic v2:请求/响应数据契约(Contract-First),自动完成参数校验与序列化。业务场景我们换成更贴近政企资产管理实际的 "设备台账管理"(字段含资产编号、类别、状态、存放位置、采购日期),比"用户余额"更能体现国产库在真实业务系统中的用法。读完后你可以把它当作模板,平滑扩展到更复杂的联表查询、审计字段、软删除等需求。

技术选型对比(为什么不用 Flask + 裸 SQL)

维度

传统示例(Flask + 裸 psycopg2)

本文方案(FastAPI + SQLAlchemy 2.0)

参数校验

手写 request.json 判断,易漏

Pydantic 自动校验,非法即 422

SQL 安全

依赖开发者自觉用 %s 占位

ORM 参数化,从根本上防注入

可维护性

SQL 散落在路由里

模型/仓储/路由分层,关注点分离

文档

需额外写

/docs 自动生成 OpenAPI

并发模型

阻塞 WSGI

异步 ASGI,连接池复用更高效

一、环境准备磨刀不误砍柴工,先把运行地基打好。1.1 前置条件

组件

要求

说明

KingbaseES

V8R6 / V9 及以上

默认监听端口 54321(非 PostgreSQL 的 5432)

Python

3.9+

推荐 3.11,类型标注与运行性能更优

数据库账号

具备目标库建表/DML 权限

演示沿用默认 SYSTEM 用户

常见踩坑:KES 默认端口是 54321,很多人照抄 PG 的 5432 连半天超时;Linux 上若报 libkci.so: cannot open shared object file,需把 KES 安装目录的 lib 路径加入 LD_LIBRARY_PATH(如 export LD_LIBRARY_PATH=/opt/Kingbase/ES/V9/lib:$LD_LIBRARY_PATH)。1.2 安装依赖金仓官方驱动 ksycopg2 已发布到 PyPI(x86_64 / Windows 64 位可直接 pip 安装)。个人习惯用 uv 管理依赖,用 pip 也完全等价:

代码语言:txt
复制
# 方式 A:uv(推荐)
uv add fastapi uvicorn "sqlalchemy>=2.0" pydantic-settings ksycopg2

# 方式 B:pip
pip install "fastapi>=0.110" "uvicorn[standard]" "sqlalchemy>=2.0" \
            "pydantic-settings>=2.0" ksycopg2

若 pip install ksycopg2 在你的平台拉不到对应架构的包(如鲲鹏、龙芯等信创芯片),前往金仓官网下载页取对应平台的驱动压缩包,解压后把 ksycopg2 目录放入 site-packages 即可。二、项目结构设计我们不做单文件堆代码,而是按分层架构(Layered Architecture)组织,方便后续接入单测、CI 与容器化部署:

分层的价值:路由层只负责 HTTP 语义(状态码、报文),业务逻辑下沉到 repositories。哪天要换 Web 框架(如 Django)或加 Redis 缓存,改动被锁死在最小影响面,符合"关注点分离"原则。三、完整代码实现与详解下面按"配置 → 连接 → 建表 → CRUD → 启动"的顺序逐块拆解,每段均可直接复制到对应文件运行。3.1 依赖声明与项目根配置先用 pyproject.toml 把依赖与 Python 版本锁定,保证可复现构建:

代码语言:txt
复制
[project]
name = "kingbase-device-api"
version = "1.0.0"
requires-python = ">=3.9"
dependencies = [
    "fastapi>=0.110",
    "uvicorn[standard]>=0.29",
    "sqlalchemy>=2.0",
    "pydantic-settings>=2.0",
    "ksycopg2>=1.0",
]

[tool.uv]
package = false

.env.example(明文密钥绝不进版本库):

代码语言:txt
复制
KES_DB_HOST=127.0.0.1
KES_DB_PORT=54321
KES_DB_USER=SYSTEM
KES_DB_PASSWORD=your_strong_password
KES_DB_NAME=TEST

3.2 配置中心 config.py所有可变配置收口到一处,用 pydantic-settings 读取环境变量,杜绝密码硬编码:

代码语言:txt
复制
# src/app/config.py
from functools import lru_cache

from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
    """应用配置,统一从环境变量 / .env 读取。"""

    model_config = SettingsConfigDict(
        env_file=".env", env_prefix="KES_", extra="ignore"
    )

    db_host: str = "127.0.0.1"
    db_port: int = 54321
    db_user: str = "SYSTEM"
    db_password: str = "123456"
    db_name: str = "TEST"
    app_name: str = "kingbase-device-api"
    api_prefix: str = "/api/v1"

    @property
    def database_url(self) -> str:
        """拼接 SQLAlchemy 连接串。

        KingbaseES 的 SQLAlchemy dialect 名为 ``kingbase``,
        官方驱动 ``ksycopg2`` 为其默认 driver,故可省略不写。
        """
        return (
            f"kingbase+ksycopg2://{self.db_user}:{self.db_password}"
            f"@{self.db_host}:{self.db_port}/{self.db_name}"
        )


@lru_cache
def get_settings() -> Settings:
    """返回进程级单例配置,避免每次请求重复实例化。"""
    return Settings()

3.3 数据库引擎与连接池 database.py这是"连接"的核心。我们用连接池复用物理连接,而非每次请求新建,显著降低握手开销:

代码语言:txt
复制
# src/app/database.py
from sqlalchemy import create_engine
from sqlalchemy.orm import DeclarativeBase, sessionmaker

from app.config import get_settings

settings = get_settings()

# pool_pre_ping=True:借出连接前先发轻量探活包,
# 自动剔除被服务端回收的"半开死连接",规避偶发 operational error。
engine = create_engine(
    settings.database_url,
    pool_pre_ping=True,
    pool_size=5,
    max_overflow=10,
    future=True,
)

# sessionmaker 产出数据库会话;autoflush=False 让提交时机更可控。
SessionLocal = sessionmaker(
    bind=engine, autoflush=False, autocommit=False, future=True
)


class Base(DeclarativeBase):
    """所有 ORM 模型的声明基类。"""


def get_db():
    """FastAPI 依赖:每请求独享一个 Session,用毕必关。

    yield 写法保证即便视图抛异常,finally 仍会执行关闭,
    杜绝数据库连接泄漏。
    """
    db = SessionLocal()
    try:
        yield db
    finally:
        db.close()

3.4 定义 ORM 模型 models/device.py以 Python 类声明"设备台账表",字段类型与约束一目了然:

代码语言:txt
复制
# src/app/models/device.py
from datetime import date, datetime

from sqlalchemy import Date, DateTime, Integer, String, func
from sqlalchemy.orm import Mapped, mapped_column

from app.database import Base


class Device(Base):
    """设备台账表 device_asset。"""

    __tablename__ = "device_asset"

    id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
    # 资产编号唯一,避免重复登记;建索引加速按编号检索
    asset_code: Mapped[str] = mapped_column(String(32), unique=True, nullable=False, index=True)
    name: Mapped[str] = mapped_column(String(128), nullable=False)
    category: Mapped[str] = mapped_column(String(64), nullable=False)
    # 状态给默认值"在用",降低入库必填负担
    status: Mapped[str] = mapped_column(String(16), nullable=False, default="在用")
    location: Mapped[str] = mapped_column(String(128), nullable=False)
    purchase_date: Mapped[date] = mapped_column(Date, nullable=False)
    created_at: Mapped[datetime] = mapped_column(DateTime, server_default=func.now())
    # onupdate 使每次 UPDATE 自动刷新时间戳,无需业务代码手动维护
    updated_at: Mapped[datetime] = mapped_column(
        DateTime, server_default=func.now(), onupdate=func.now()
    )

本方案的"建表"不再需要单独的 HTTP 接口——应用在启动时由 Base.metadata.create_all 自动完成(见 3.10)。生产环境建议改用语义化迁移工具 Alembic 管理表结构演进,避免直接 DDL 误操作且支持回滚。3.5 定义接口契约 schemas/device.py以 Pydantic 把"入参、出参"钉死,契约优先(Contract-First):

代码语言:txt
复制
# src/app/schemas/device.py
from datetime import date, datetime
from typing import Literal

from pydantic import BaseModel, ConfigDict, Field

# 用 Literal 约束状态枚举,非法值直接被 422 拒绝
StatusLiteral = Literal["在用", "闲置", "维修", "报废"]


class DeviceCreate(BaseModel):
    """创建设备的入参模型。"""

    asset_code: str = Field(..., max_length=32, description="资产编号,唯一")
    name: str = Field(..., max_length=128)
    category: str = Field(..., max_length=64)
    status: StatusLiteral = "在用"
    location: str = Field(..., max_length=128)
    purchase_date: date


class DeviceUpdate(BaseModel):
    """更新模型:全字段可选,仅传需变更项(PATCH 语义)。"""

    name: str | None = Field(None, max_length=128)
    category: str | None = Field(None, max_length=64)
    status: StatusLiteral | None = None
    location: str | None = Field(None, max_length=128)
    purchase_date: date | None = None


class DeviceOut(BaseModel):
    """响应模型。from_attributes=True 使 ORM 对象可直接序列化。"""

    model_config = ConfigDict(from_attributes=True)

    id: int
    asset_code: str
    name: str
    category: str
    status: str
    location: str
    purchase_date: date
    created_at: datetime
    updated_at: datetime

3.6 数据访问层 repositories/device.py将数据库操作封装为纯函数,路由层不直接持有 Session,便于单测与替换:

代码语言:txt
复制
# src/app/repositories/device.py
from typing import Optional, Tuple

from sqlalchemy import func, select
from sqlalchemy.orm import Session

from app.models.device import Device
from app.schemas.device import DeviceCreate, DeviceUpdate


def create_device(db: Session, payload: DeviceCreate) -> Device:
    """插入一条设备记录并返回完整对象。"""
    device = Device(**payload.model_dump())
    db.add(device)
    db.commit()
    db.refresh(device)  # 取回自增 id 与数据库侧默认时间戳
    return device


def get_device(db: Session, device_id: int) -> Optional[Device]:
    """按主键取单条;不存在返回 None。"""
    return db.get(Device, device_id)


def list_devices(
    db: Session, *, page: int = 1, size: int = 20, status: Optional[str] = None
) -> Tuple[list[Device], int]:
    """分页列出设备,可选按状态过滤。返回 (数据列表, 总数)。"""
    stmt = select(Device)
    count_stmt = select(func.count()).select_from(Device)
    if status:
        stmt = stmt.where(Device.status == status)
        count_stmt = count_stmt.where(Device.status == status)

    total = db.scalar(count_stmt) or 0
    # 倒序让最新录入置顶;offset/limit 实现游标式分页
    items = db.scalars(
        stmt.order_by(Device.id.desc()).offset((page - 1) * size).limit(size)
    ).all()
    return items, total


def update_device(db: Session, device_id: int, payload: DeviceUpdate) -> Optional[Device]:
    """局部更新;不存在返回 None。"""
    device = db.get(Device, device_id)
    if not device:
        return None
    # exclude_unset=True:仅覆盖前端真实传参的字段,避免把未传字段误写 None
    for key, value in payload.model_dump(exclude_unset=True).items():
        setattr(device, key, value)
    db.commit()
    db.refresh(device)
    return device


def delete_device(db: Session, device_id: int) -> bool:
    """删除设备;返回是否真正命中并删除。"""
    device = db.get(Device, device_id)
    if not device:
        return False
    db.delete(device)
    db.commit()
    return True

3.7 创建设备接口 POST /api/v1/devices路由层保持"薄"——仅做参数接收、调用仓储、包装响应:

代码语言:txt
复制
# src/app/api/v1/devices.py(片段一:创建)
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session

from app.database import get_db
from app.repositories.device import create_device
from app.schemas.device import DeviceCreate, DeviceOut

router = APIRouter(prefix="/devices", tags=["devices"])


@router.post(
    "",
    response_model=DeviceOut,
    status_code=status.HTTP_201_CREATED,
    summary="登记一台新设备",
)
def add_device(payload: DeviceCreate, db: Session = Depends(get_db)) -> DeviceOut:
    try:
        return create_device(db, payload)
    except IntegrityError:
        # asset_code 唯一约束冲突 -> 409 冲突,而非 500 内部错误
        db.rollback()
        raise HTTPException(
            status_code=status.HTTP_409_CONFLICT, detail="资产编号已存在"
        )

3.8 查询单条接口 GET /api/v1/devices/{id}

代码语言:txt
复制
# src/app/api/v1/devices.py(片段二:查询单条)
from app.repositories.device import get_device

@router.get("/{device_id}", response_model=DeviceOut, summary="按 ID 查询设备")
def read_device(device_id: int, db: Session = Depends(get_db)) -> DeviceOut:
    device = get_device(db, device_id)
    if not device:
        raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="设备不存在")
    return device

3.9 更新 / 删除 / 分页列表接口

列表接口统一返回 {total, page, size, items},前端分页组件无需另行统计总数。3.10 应用装配 main.py 与生命周期

最后记得补齐各目录的 __init__.py,并在包根建立 src/app/__init__.py,使 app 成为可导入的包。

代码语言:txt
复制
# src/app/api/v1/devices.py(片段三:更新、删除、列表)
from typing import Optional

from fastapi import Query

from app.repositories.device import delete_device, list_devices, update_device
from app.schemas.device import DeviceOut, DeviceUpdate


@router.patch("/{device_id}", response_model=DeviceOut, summary="局部更新设备")
def edit_device(
    device_id: int, payload: DeviceUpdate, db: Session = Depends(get_db)
) -> DeviceOut:
    device = update_device(db, device_id, payload)
    if not device:
        raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="设备不存在")
    return device


@router.delete("/{device_id}", status_code=status.HTTP_204_NO_CONTENT, summary="删除设备")
def remove_device(device_id: int, db: Session = Depends(get_db)) -> None:
    if not delete_device(db, device_id):
        raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="设备不存在")


@router.get("", summary="分页列出设备")
def list_device_api(
    page: int = Query(1, ge=1, description="页码"),
    size: int = Query(20, ge=1, le=100, description="每页条数"),
    status_filter: Optional[str] = Query(None, alias="status", description="按状态过滤"),
    db: Session = Depends(get_db),
) -> dict:
    items, total = list_devices(db, page=page, size=size, status=status_filter)
    return {
        "total": total,
        "page": page,
        "size": size,
        "items": [DeviceOut.model_validate(d) for d in items],
    }

四、接口测试方法启动服务(生产建议用 uvicorn 多进程,而非 Flask 的 debug=True 单进程):

启动后访问 http://127.0.0.1:8000/docs 即可看到自动生成的 Swagger 交互文档。下面用 curl 走一遍完整链路:

自动化回归可用 fastapi.testclient.TestClient 直连内存路由,配合 pytest 做断言,比手敲 curl 更稳健、可纳入 CI。

五、生产环境注意事项(避坑指南)关闭开发态自动重载:上线用 uvicorn app.main:app --workers 4,不挂 --reload;或采用 Gunicorn + Uvicorn Worker 模式。密钥必须走环境变量:.env 加入 .gitignore,任何明文密码、Token 不进版本库。连接池需按并发调优:pool_size / max_overflow 并非越大越好,须对照 KES 侧 max_connections 设定,避免压垮数据库。SQL 注入免疫:ORM 参数化查询天然防注入;凡手写明文字符串 SQL,一律用 text() + 命名参数绑定,绝不拼接变量。健康检查与可观测性:补充 GET /healthz 探活接口,接入 Prometheus 监控连接池占用;异常统一经中间件转为结构化日志。鉴权不可省略:示例无认证,生产必须落地 API Key / JWT / OAuth2;KES 侧建议启用 scram-sha-256,跨网络传输开启 SSL(sslmode=verify-ca)。事务边界清晰:写操作确保 commit/rollback 成对;批量导入用 bulk_save_objects 并控制单批规模,避免长事务锁表。迁移用 Alembic:create_all 仅适合 Demo,正式项目用迁移脚本管理表结构演进,可回滚、可审计。隐藏真实出口 IP:当 API 需对公网或第三方暴露时,建议前置反向代理/网关(如 Nginx、APISIX),将真实服务与数据库 IP 收敛在内网,降低被直连扫描与溯源的风险。

六、实战扩展:当 API 成为数据采集后端(亿牛云代理 IP 选型)在真实的政企数据平台中,上面这套 KES API 常常处在数据采集链路的最末端——上游的爬虫/采集程序把抓取到的设备、资产或舆情数据回写进库。一旦采集端需要高频访问外部目标站点,就会面临两个工程问题:反爬封禁与源站 IP 暴露。6.1 为什么需要代理 IP 层IP 轮换规避封禁:目标站点对单一出口 IP 限频/拉黑时,代理池自动切换出口,保障采集连续性;隐藏真实源站:采集服务器与后端 API 的真实公网 IP 不对外暴露,降低被溯源与针对性攻击的风险;企业级稳定可用:相比自建代理,专业服务商提供 SLA 保障、自动鉴权与高可用隧道,运维成本更低。在代理 IP 选型上,亿牛云(企业级代理 IP 服务商)的隧道代理是这类场景的常见方案:它提供稳定的 HTTP/HTTPS 隧道入口,客户端只需配置一个代理地址,出口 IP 由亿牛云侧自动轮换,无需自建与维护 IP 池。6.2 采集端经亿牛云代理回写 KES API下面给出采集端(异步 httpx)通过亿牛云隧道代理,把抓取结果写入我们上面 FastAPI 接口的示例:

代码语言:txt
复制
# collector.py —— 采集端示例(与 KES API 解耦的独立进程)
import httpx

# 亿牛云隧道代理入口(企业级,出口 IP 自动轮换)
# 格式:http://<用户名>:<密码>@<网关地址>:<端口>
YINIU_PROXY = "http://your_user:your_pass@gateway.ipyiniu.com:埠"

API_ENDPOINT = "http://api.your-intranet.com/api/v1/devices"


async def push_device(record: dict) -> int:
    """将一条采集到的设备信息经亿牛云代理回写至 KES API。

    Args:
        record: 待写入的设备字典,字段需符合 DeviceCreate 契约

    Returns:
        HTTP 状态码(201 表示创建成功)

    Raises:
        httpx.HTTPStatusError: 服务端返回 4xx/5xx 时抛出
    """
    # proxy 参数让本次请求走亿牛云隧道,出口 IP 由服务端轮换
    async with httpx.AsyncClient(proxy=YINIU_PROXY, timeout=10.0) as client:
        resp = await client.post(API_ENDPOINT, json=record)
        resp.raise_for_status()
        return resp.status_code


# 典型采集循环(伪代码)
# for item in crawl_target_site():
#     await push_device(item)   # 经亿牛云代理写入金仓

6.3 选型与合规要点隧道 vs 短效代理:隧道代理适合长稳采集、免管理 IP 池;短效/动态代理适合需要精细控制单 IP 时长的场景。鉴权与限速:亿牛云等合规服务商均要求账号鉴权,务必将密钥放入环境变量,并按目标站点 robots.txt 与法律法规控制采集频率。链路分工:代理 IP 负责"采集端出口",KES API 与数据库仍应置于内网;公网只暴露经鉴权的 API 网关,形成"采集侧代理 + 内网存储"的清晰边界。至此,一条"亿牛云代理采集 → FastAPI 接收 → KingbaseES 落库"的国产化数据链路就闭环了。七、总结本文以 FastAPI + SQLAlchemy 2.0 + 金仓官方驱动 ksycopg2 落地了一套可直接上线的 KES 数据 API:配置中心化、连接池化、模型 ORM 化、接口契约化。得益于 KingbaseES 对 PostgreSQL 协议的兼容,整套代码几乎可零修改地理解为"在连一台 PG"——这正是国产数据库生态成熟的红利,也是信创替代中迁移成本可控的关键。你已经掌握了"连库 → 建表 → CRUD → 暴露 API"的完整闭环。进一步的演进方向包括:补充软删除字段 deleted_at、做联表查询(如设备归属部门)、引入 Redis 缓存热点列表、用 Alembic 接管表结构版本,以及在采集侧叠加亿牛云代理 IP 构建稳定的数据入口。国产数据库的 Python 工程化之路,至此打通。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档