Skip to main content

02. 数据库与持久化:PostgreSQL、SQL、ORM

读完这篇你能:

  • 用 Docker 起一个 PostgreSQL,GUI 工具能连上
  • 会写基础 SQL(CRUD、JOIN、聚合)
  • 用 SQLModel 把 Postgres 接到 FastAPI,数据能存能取
  • 知道事务、索引是干嘛的,会用

对应架构层:数据层(核心)

前置知识:01. Web 基础

预计用时:3 小时(含动手)


0. 为什么算法工程师要学这个

你的 LLM 应用跑起来,跑了 10 分钟,服务器重启——所有用户的对话历史、积分、反馈全没了

为什么?因为你的数据存在内存(Python 变量)里。内存的特点:

  • 断电即失
  • 进程重启即失
  • 多个进程之间不共享

数据库就是解决这个问题的:把数据持久化到磁盘,即使服务器被雷劈了,只要磁盘没坏,数据就还在。

算法工程师可能会想:"我以前都用 CSV / JSON / Parquet 存数据,够用啊?"

生产环境不行,因为:

场景CSV/JSON数据库
100 个用户同时写文件锁,排队,慢死并发安全
按条件查询(user_id=123 的订单)全表扫描,O(N)用索引,O(log N)
写一半断电文件损坏ACID 事务保证
两个用户同时改同一条后写的覆盖先写的行锁、事务隔离
数据量到千万行打开就卡死索引、分片

结论:任何"需要长期保存"的数据,都要用数据库。

这一篇我们用 PostgreSQL——开源数据库的事实标准,功能强大(支持 JSON、向量、全文检索)、社区活跃、主流大厂都在用。


1. 核心概念(最少必要理论)

1.1 关系型数据库 vs NoSQL

类型例子特点什么时候用
关系型(RDBMS)PostgreSQL、MySQL、SQLite表结构固定、SQL、强一致、事务业务数据(用户、订单)
键值(KV)Redis、Memcached超快、内存为主缓存、session
文档型MongoDB、CouchDBschema 灵活、JSON 存储非结构化、原型快速
列族Cassandra、HBase海量写入、列查询时序、日志
Neo4j关系密集(社交、推荐)知识图谱
向量pgvector、Qdrant、Milvus向量相似度检索RAG、语义搜索

本篇讲关系型,其他类型在 03. 缓存与存储 里讲。

为什么关系型优先?

因为业务数据 90% 是结构化的(用户有 id、name、email、created_at),关系型用 40 年证明了它最稳。

1.2 表、行、列

数据库 → 表(Table)→ 行(Row,一条记录)→ 列(Column,一个字段)。

类比 Pandas:数据库 = 一堆 DataFrame,表 = 一个 DataFrame,行 = iloc[i],列 = df['name']

users 表
┌────┬───────────┬─────────────────┬────────────┐
│ id │ name │ email │ created_at │
├────┼───────────┼─────────────────┼────────────┤
│ 1 │ sanbu │ [email protected] │ 2026-06-23 │
│ 2 │ alice │ [email protected] │ 2026-06-22 │
│ 3 │ bob │ [email protected] │ 2026-06-21 │
└────┴───────────┴─────────────────┴────────────┘
  • id主键(Primary Key),唯一标识一行
  • email 可以加唯一约束(UNIQUE),不允许重复

1.3 SQL 四大类(必背)

CRUD = Create / Read / Update / Delete:

操作SQL 语句
创建INSERT INTO users (name, email) VALUES ('sanbu', '[email protected]')
查询SELECT * FROM users WHERE id = 1
更新UPDATE users SET name = 'new' WHERE id = 1
删除DELETE FROM users WHERE id = 1

记住:永远带 WHEREDELETE FROM users(没 where)会把全表删光。

1.4 JOIN(连接多表)

为什么拆成多张表?避免冗余

如果每个订单都把用户信息复制一遍,改个邮箱要改 100 万条订单。拆开后:

  • users:id, name, email
  • orders:id, user_id, amount, created_at

orders.user_id外键(Foreign Key),指向 users.id

要查询"alice 买过什么",用 JOIN:

SELECT users.name, orders.amount, orders.created_at
FROM users
JOIN orders ON users.id = orders.user_id
WHERE users.name = 'alice';

返回:

name  | amount | created_at
------+--------+---------------------
alice | 99.00 | 2026-06-22 10:00:00
alice | 15.50 | 2026-06-22 11:30:00

1.5 索引(Index)——算法工程师最容易理解的优化

问题:表有 1000 万行,SELECT * FROM users WHERE email = '[email protected]' 要扫 1000 万次?

解决:给 email 加索引。Postgres 自动维护一个 B-Tree,查询变成 O(log N)。

类比:字典的目录。没有目录,找一个字要从头翻到尾(O(N));有目录,先翻目录到 P 开头(O(log N)),再翻到那一页(O(1))。

CREATE INDEX idx_users_email ON users(email);

索引的代价:

  • 占磁盘空间(索引本身也是数据)
  • 写入变慢(每写一行,索引也要更新)
  • 不是越多越好——只给查询频繁的列加

什么时候该加索引:

  • WHERE 频繁过滤的列
  • JOIN 的连接列(user_id)
  • 排序列(ORDER BY created_at)

1.6 事务(ACID)

事务(Transaction) 是"一组操作,要么全成功,要么全失败"。

经典例子:转账。

BEGIN;
UPDATE accounts SET balance = balance - 100 WHERE user_id = 1; -- 扣钱
UPDATE accounts SET balance = balance + 100 WHERE user_id = 2; -- 加钱
COMMIT; -- 或者 ROLLBACK 撤销

如果第二条失败,第一条也得撤销,不然用户 1 钱没了,用户 2 没收到。

ACID 四特性:

  • Atomicity(原子性):全成功或全失败
  • Consistency(一致性):数据始终满足约束(如余额不能为负)
  • Isolation(隔离性):并发事务互不干扰
  • Durability(持久性):提交后断电也不丢

算法工程师类比:事务像"训练 checkpoint 的原子写"——要么 checkpoint 完整保存,要么不存,不能存一半。

1.7 ORM:用 Python 对象操作数据库

直接写 SQL 字符串容易出错:

# ❌ 容易拼错、SQL 注入风险
cursor.execute(f"SELECT * FROM users WHERE name = '{user_input}'")

ORM(Object-Relational Mapping)把表映射成 Python 类:

# 表 → 类
class User(SQLModel, table=True):
id: int | None = Field(default=None, primary_key=True)
name: str
email: str

# 行 → 实例
user = User(name="sanbu", email="[email protected]")

# 查询 → Pythonic API
session.exec(select(User).where(User.name == "sanbu")).all()

好处:

  • 类型安全(IDE 有补全)
  • 防 SQL 注入(自动参数化)
  • 跨数据库(切 MySQL/SQLite 几乎不改代码)

本系列用 SQLModel(FastAPI 作者写的,结合了 SQLAlchemy 和 Pydantic)。其他选择:SQLAlchemy(老牌)、Django ORM、Tortoise ORM。


2. 实战:给 LLM 应用建数据库

目标:给上一篇的 LLM 应用加 3 张表——usersqueriesfeedback,完整 CRUD。

2.1 用 Docker 起 PostgreSQL

为什么用 Docker?——不用装 Postgres,不污染系统,一行命令起一个。

创建 docker-compose.yml:

services:
postgres:
image: postgres:16
environment:
POSTGRES_USER: myuser
POSTGRES_PASSWORD: mypassword
POSTGRES_DB: myapp
ports:
- "5432:5432"
volumes:
- pgdata:/var/lib/postgresql/data

volumes:
pgdata:

启动:

docker compose up -d
docker compose ps # 确认在跑

连数据库(GUI 推荐):

  • VS Code 装 Database Client 插件(cweijan)
  • 或下载 DBeaver / TablePlus

连接参数:host=localhost, port=5432, user=myuser, password=mypassword, db=myapp。

命令行:

docker compose exec postgres psql -U myuser -d myapp

2.2 定义模型

新建 models.py:

from datetime import datetime
from typing import Optional
from sqlmodel import SQLModel, Field

class User(SQLModel, table=True):
__tablename__ = "users"

id: Optional[int] = Field(default=None, primary_key=True)
email: str = Field(unique=True, index=True)
name: str
hashed_password: str # 04 篇讲怎么 hash
created_at: datetime = Field(default_factory=datetime.utcnow)

class Query(SQLModel, table=True):
__tablename__ = "queries"

id: Optional[int] = Field(default=None, primary_key=True)
user_id: int = Field(foreign_key="users.id", index=True)
input_text: str
output_text: str
tokens_used: int
latency_ms: int
created_at: datetime = Field(default_factory=datetime.utcnow, index=True)

class Feedback(SQLModel, table=True):
__tablename__ = "feedback"

id: Optional[int] = Field(default=None, primary_key=True)
query_id: int = Field(foreign_key="queries.id", index=True)
user_id: int = Field(foreign_key="users.id")
rating: int # 1-5
comment: Optional[str] = None
created_at: datetime = Field(default_factory=datetime.utcnow)

关键点:

  • table=True 表示这是个表,不是普通 Pydantic 模型
  • Field(primary_key=True) 主键
  • Field(foreign_key="users.id") 外键
  • Field(index=True) 自动建索引
  • Optional[int] 表示可为 None(自增主键)

2.3 数据库连接和 Session

新建 db.py:

from sqlmodel import create_engine, Session

DATABASE_URL = "postgresql://myuser:mypassword@localhost:5432/myapp"

engine = create_engine(DATABASE_URL, echo=True) # echo=True 打印 SQL,调试用

def get_session():
with Session(engine) as session:
yield session

生产环境:

  • echo=False(关掉 SQL 日志)
  • 加连接池参数:create_engine(DATABASE_URL, pool_size=10, max_overflow=20)
  • 异步用 create_async_engine(本篇先用同步,异步见 04 篇)

2.4 建表

新建 init_db.py:

from sqlmodel import SQLModel
from db import engine
from models import User, Query, Feedback

def init_db():
SQLModel.metadata.create_all(engine)

if __name__ == "__main__":
init_db()
print("Tables created")

跑一次:

python init_db.py

GUI 里刷新,看到三张表。

注意:create_all 是开发期用的,生产用 Alembic 做迁移(见 3.4)。

2.5 把数据库接进 FastAPI

main.py:

from fastapi import FastAPI, Depends, HTTPException
from sqlmodel import Session, select
from db import get_session
from models import User, Query, Feedback

app = FastAPI()

# ====== User Routes ======
@app.post("/api/users", response_model=User)
def create_user(user: User, session: Session = Depends(get_session)):
session.add(user)
session.commit()
session.refresh(user)
return user

@app.get("/api/users/{user_id}", response_model=User)
def get_user(user_id: int, session: Session = Depends(get_session)):
user = session.get(User, user_id)
if not user:
raise HTTPException(404, "User not found")
return user

@app.get("/api/users", response_model=list[User])
def list_users(limit: int = 20, offset: int = 0, session: Session = Depends(get_session)):
return session.exec(select(User).limit(limit).offset(offset)).all()

关键点:

  • Depends(get_session) 是 FastAPI 的依赖注入——每个请求自动拿到一个 session,请求结束自动关闭
  • session.add() + session.commit() 才会真正写入
  • session.refresh(user) 重新加载,拿到数据库生成的字段(如自增 id)
  • session.get(Model, pk) 按主键查
  • select(Model).where(...) 复杂查询

2.6 把 summarize 接口改成持久化

把 01 篇的 /api/summarize 改成:每次调用都记到 queries 表。

from datetime import datetime
import time

class SummarizeRequest(BaseModel):
text: str = Field(..., min_length=10)
user_id: int # 先手动传,04 篇用 JWT 自动取
max_length: int = 200

@app.post("/api/summarize")
async def summarize(req: SummarizeRequest, session: Session = Depends(get_session)):
user = session.get(User, req.user_id)
if not user:
raise HTTPException(404, "User not found")

start = time.time()
resp = await client.chat.completions.create(
model="gpt-4o-mini",
messages=[
{"role": "system", "content": "总结以下文本"},
{"role": "user", "content": req.text},
],
)
latency_ms = int((time.time() - start) * 1000)

# 存到数据库
query = Query(
user_id=req.user_id,
input_text=req.text,
output_text=resp.choices[0].message.content,
tokens_used=resp.usage.total_tokens,
latency_ms=latency_ms,
)
session.add(query)
session.commit()
session.refresh(query)

return {"summary": query.output_text, "query_id": query.id}

测试:

# 先建用户
curl -X POST http://localhost:8000/api/users \
-H "Content-Type: application/json" \
-d '{"email":"[email protected]","name":"sanbu","hashed_password":"xxx"}'

# 再调 summarize(假设 user_id=1)
curl -X POST http://localhost:8000/api/summarize \
-H "Content-Type: application/json" \
-d '{"text":"FastAPI 是一个现代的 Python Web 框架...","user_id":1}'

GUI 里看 queries 表,多了一条记录!

2.7 查询统计:JOIN + 聚合

需求:查询"每个用户的调用次数和总 token 数"。

直接在 endpoint 里写 SQL:

from sqlmodel import func

@app.get("/api/stats/users")
def user_stats(session: Session = Depends(get_session)):
stmt = (
select(
User.name,
func.count(Query.id).label("query_count"),
func.sum(Query.tokens_used).label("total_tokens"),
)
.join(Query, User.id == Query.user_id)
.group_by(User.id)
)
return session.exec(stmt).all()

等价 SQL:

SELECT users.name,
COUNT(queries.id) AS query_count,
SUM(queries.tokens_used) AS total_tokens
FROM users
JOIN queries ON users.id = queries.user_id
GROUP BY users.id;

返回:

[
{"name": "sanbu", "query_count": 42, "total_tokens": 12345},
{"name": "alice", "query_count": 13, "total_tokens": 4321}
]

3. 进阶:生产级改造

3.1 迁移工具 Alembic(必学)

问题:你已经上线了,users 表有 100 万行数据。现在想加一个字段 phone

不能用 create_all(那只建新表,不改已有的)。直接 ALTER TABLE 风险高。

Alembic 解决"数据库 schema 演进"问题——把每次改表写成脚本,可以正向 / 反向执行。

pip install alembic
alembic init alembic

alembic.ini 里的 sqlalchemy.url,改 alembic/env.py 让它认识你的模型。

创建迁移:

alembic revision --autogenerate -m "add phone to users"

会在 alembic/versions/xxx_add_phone_to_users.py 生成脚本,你检查一下:

def upgrade():
op.add_column('users', sa.Column('phone', sa.String(), nullable=True))

def downgrade():
op.drop_column('users', 'phone')

执行:

alembic upgrade head    # 升到最新
alembic downgrade -1 # 回退一版

生产流程:

  1. 改模型(加字段)
  2. alembic revision --autogenerate
  3. 人工检查生成的脚本(autogenerate 不完美)
  4. 在 staging 跑一遍,验证
  5. 上线时 alembic upgrade head

3.2 索引优化:EXPLAIN

查询慢?用 EXPLAIN 看执行计划:

EXPLAIN ANALYZE SELECT * FROM queries WHERE user_id = 1;

如果看到 Seq Scan(Sequential Scan,全表扫描),说明没命中索引。 如果看到 Index Scan,命中了。

复合索引:如果经常按 (user_id, created_at) 一起查,建复合索引:

CREATE INDEX idx_queries_user_created ON queries(user_id, created_at);

3.3 连接池

默认 create_engine 已经有连接池(pool_size=5)。生产环境要调:

engine = create_engine(
DATABASE_URL,
pool_size=20, # 常驻连接数
max_overflow=10, # 突发可借的额外连接
pool_timeout=30, # 等连接的超时
pool_recycle=1800, # 连接最长寿命(秒),避免数据库主动断开
)

经验值:

  • CPU 密集型:pool_size = CPU 核数 × 2 + 1
  • IO 密集型(LLM 应用):pool_size 可以更大,20-50
  • 单个 Postgres 数据库默认 max_connections=100,所有服务加起来不能超

3.4 异步 SQL(配合 async FastAPI)

01 篇讲了 LLM 调用要 async。数据库同样——如果你用同步 SQLAlchemy,在 async 路由里会阻塞事件循环。

方案:用 aiosqliteasyncpgdatabases 等异步驱动。SQLModel 1.0+ 支持异步:

from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession

engine = create_async_engine("postgresql+asyncpg://user:pass@host/db")

async def get_session():
async with AsyncSession(engine) as session:
yield session

或者偷懒:数据库操作用同步(LLM 调用用 async),让 FastAPI 把同步路由放到线程池跑:

@app.post("/api/summarize")
def summarize(req, session = Depends(get_session)): # 注意:同步 def
# 同步 SQL,但 LLM 调用用 asyncio.run(...) 包一下

这种方式简单但有性能损失。推荐全异步

3.5 备份和恢复

手动备份:

# 备份
docker compose exec postgres pg_dump -U myuser myapp > backup.sql

# 恢复
cat backup.sql | docker compose exec -T postgres psql -U myuser myapp

自动备份:生产环境用云数据库(RDS、Cloud SQL),自动每天备份 + 增量备份。绝对不要让自己管 Postgres 备份。


4. 常见踩坑

坑 1:N+1 查询

# ❌ 取 100 个用户,每人查一次最近订单 = 101 次 SQL
users = session.exec(select(User)).all()
for u in users:
u.last_order = session.exec(
select(Order).where(Order.user_id == u.id).order_by(Order.id.desc()).limit(1)
).first()

症状:页面上 100 个用户加载要 5 秒。

解决:用 JOIN 一次查完,或用 selectinload 预加载关系。

坑 2:忘了 commit

session.add(user)
# 忘了 session.commit()
return user # 返回的对象有 id,但其实没入库

坑 3:把密码明文存数据库

# ❌ 数据库泄露 = 所有用户密码泄露
user.hashed_password = "plaintext123"

# ✅ 用 bcrypt / argon2
from passlib.context import CryptContext
pwd_ctx = CryptContext(schemes=["bcrypt"])
user.hashed_password = pwd_ctx.hash(password)

详见 04. Web 服务

坑 4:float 存金额

# ❌ 浮点精度问题,0.1 + 0.2 = 0.30000000000000004
amount: float

# ✅ 用 Numeric / Decimal
from decimal import Decimal
amount: Decimal = Field(sa_type=Numeric(10, 2))

坑 5:不处理并发写冲突

两个请求同时给用户加 100 积分:

# ❌ 读改写,并发会丢更新
user = session.get(User, 1)
user.credits += 100
session.commit()

解决:用原子 UPDATE:

# ✅ 一条 SQL 完成
session.exec(
update(User).where(User.id == 1).values(credits=User.credits + 100)
)

或者用乐观锁(version 字段)+ 重试。


5. 检查清单

  • 我能用 Docker 起一个 Postgres,GUI 能连上
  • 我会写基础 SQL:INSERT、SELECT、UPDATE、DELETE、JOIN、GROUP BY
  • 我知道主键、外键、唯一约束、索引分别是什么
  • 我能解释 ACID 四特性
  • 我会用 SQLModel 定义模型,FastAPI 能 CRUD
  • 我会用 Depends(get_session) 做依赖注入
  • 我知道 N+1 查询是什么,会避免
  • 我会用 Alembic 做迁移
  • 我会看 EXPLAIN 判断有没有命中索引
  • 我知道密码不能明文存,要 hash

6. 下一步

你已经能给应用加"长期记忆"了。下一篇继续强化数据层:

学完 03,数据层就完整了。然后 04 篇回到应用层——加认证、权限、流式输出。

进阶学习资源