Python 数据库访问
大多数后端程序最终都要保存数据,例如用户、订单、采集任务、审计日志、AI 对话记录、资产目录。Python 访问数据库不要只会 execute(),还要理解连接、游标、事务、连接池、ORM、迁移和慢 SQL 排查。
本页以关系型数据库为主,例如 MySQL、PostgreSQL、SQLite。不同数据库语法细节不同,但 Python 访问数据库的核心流程相通。
学习目标
学完本页你应该能回答:
- Python 是怎么通过驱动访问数据库的。
- Connection、Cursor、Transaction 分别是什么。
- 为什么参数绑定能防 SQL 注入。
- 事务什么时候提交、什么时候回滚。
- 连接池为什么需要,为什么不能无限大。
- SQLAlchemy Core、ORM、Session、Alembic 分别解决什么问题。
- 大查询、慢 SQL、锁等待、连接耗尽怎么排查。
- 商业项目中 Repository 和 Service 的事务边界怎么划分。
数据库访问的基本流程
flowchart TD
A["准备连接信息"] --> B["数据库驱动创建连接"]
B --> C["创建 Cursor 或 Session"]
C --> D["发送 SQL 和参数"]
D --> E["数据库解析和执行 SQL"]
E --> F["驱动接收结果集"]
F --> G["Python 读取结果"]
G --> H{"是否修改数据"}
H -- "是" --> I["commit 或 rollback"]
H -- "否" --> J["关闭或归还连接"]
I --> J核心原则:Python 不直接操作数据库文件或表结构,而是通过数据库驱动把 SQL 发给数据库服务,由数据库完成解析、优化、执行、锁控制和事务处理。
常见方案怎么选
| 方案 | 说明 | 适合场景 |
|---|---|---|
sqlite3 | Python 标准库自带,单文件数据库 | 学习、小工具、本地缓存 |
| PyMySQL / mysqlclient | MySQL 驱动 | 连接 MySQL |
| psycopg | PostgreSQL 驱动 | 连接 PostgreSQL |
| SQLAlchemy Core | 用 Python API 构造 SQL | 想减少手写 SQL,又要控制 SQL |
| SQLAlchemy ORM | 类和表映射,通过对象操作数据 | 中大型 Web 项目 |
| Alembic | 数据库迁移工具 | 管理表结构版本 |
学习路线:
- 先用
sqlite3理解连接、游标、SQL、事务。 - 再用 MySQL/PostgreSQL 驱动理解网络数据库。
- 再学 SQLAlchemy Session 生命周期和事务边界。
- 最后学 Alembic,把表结构变更纳入版本管理。
DB-API:Connection 和 Cursor
Python 数据库驱动大多遵循 DB-API 思路。
| 对象 | 作用 | 类比 |
|---|---|---|
| Connection | 到数据库的一条连接,承载事务 | 一条电话线 |
| Cursor | 执行 SQL、读取结果的游标 | 接线员 |
| Transaction | 一组要么全成功要么全失败的操作 | 一份完整业务单据 |
最小 Demo:
import sqlite3
conn = sqlite3.connect("demo.db")
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE IF NOT EXISTS users (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL,
age INTEGER NOT NULL
)
""")
cursor.execute("INSERT INTO users (name, age) VALUES (?, ?)", ("张三", 18))
conn.commit()
cursor.execute("SELECT id, name, age FROM users")
for row in cursor.fetchall():
print(row)
cursor.close()
conn.close()这段代码的执行过程:
sqlite3.connect()打开数据库连接。conn.cursor()创建游标。cursor.execute()把 SQL 发送给数据库。- 插入、更新、删除后要
commit(),否则事务可能没有真正提交。 - 查询通过
fetchone()、fetchmany()、fetchall()获取结果。 - 最后关闭游标和连接。
驱动到底做了什么
Python 代码不会直接“进入数据库表里拿数据”。真正的链路通常是:Python 调用驱动,驱动把 SQL 和参数按照数据库协议发给数据库服务,数据库执行完成后再把结果返回给驱动。
flowchart TD
A["Python 业务代码"] --> B["数据库驱动"]
B --> C["序列化 SQL 和参数"]
C --> D["通过 TCP 发送给数据库"]
D --> E["数据库认证和权限检查"]
E --> F["解析 SQL"]
F --> G["优化器选择执行计划"]
G --> H["执行器读取索引和数据页"]
H --> I["返回结果集或影响行数"]
I --> J["驱动转换成 Python 对象"]这条链路解释了几个常见现象:
- SQL 写错不是 Python 解释器发现的,而是数据库返回语法错误。
- 网络数据库连接失败可能是账号、密码、端口、防火墙、网络、TLS、权限任意一层出问题。
- 查询慢不一定是 Python 慢,通常要看数据库执行计划、锁等待、IO、索引和返回行数。
- 驱动会把数据库类型转换成 Python 类型,例如 SQL 的
VARCHAR变成str,INT变成int,时间类型变成datetime。
SQLite 比较特殊,它是进程内文件数据库,不需要 TCP 连接数据库服务;MySQL、PostgreSQL 则通常是通过网络连接数据库进程。
| 数据库 | 连接方式 | 适合 |
|---|---|---|
| SQLite | 本地文件 | 学习、本地工具、轻量缓存 |
| MySQL | TCP 连接数据库服务 | Web 业务、订单、资产、权限 |
| PostgreSQL | TCP 连接数据库服务 | 复杂查询、GIS、强 SQL 能力 |
execute、executemany 和影响行数
单条执行:
cursor.execute(
"INSERT INTO users (name, age) VALUES (?, ?)",
("张三", 18),
)批量执行:
rows = [
("张三", 18),
("李四", 20),
("王五", 22),
]
cursor.executemany("INSERT INTO users (name, age) VALUES (?, ?)", rows)executemany 的价值是减少 Python 和数据库之间的来回通信次数。但它也不是“越大越好”,一次批量太大可能导致:
- SQL 包太大。
- 事务持有锁太久。
- 失败回滚成本高。
- 内存占用升高。
商业批处理建议按批次提交,例如每 500 或 1000 行一批,具体要结合数据库性能、锁竞争和失败补偿策略。
影响行数:
cursor.execute("UPDATE users SET age = ? WHERE id = ?", (19, 1))
print(cursor.rowcount)rowcount 常用于判断更新是否真正命中数据。例如订单支付更新状态时,如果影响行数为 0,可能说明订单不存在、状态已经变化或并发抢占失败。
参数绑定为什么能防 SQL 注入
危险写法:
name = "' OR '1'='1"
sql = f"SELECT * FROM users WHERE name = '{name}'"
cursor.execute(sql)最终 SQL 可能变成:
SELECT * FROM users WHERE name = '' OR '1'='1'这会绕过原本的过滤条件。
正确写法:
cursor.execute("SELECT * FROM users WHERE name = ?", (name,))MySQL 常见参数符号:
cursor.execute("SELECT * FROM users WHERE name = %s", (name,))参数绑定的原理是:SQL 模板和参数分开发给驱动,驱动把用户输入当作“值”处理,而不是当作 SQL 语法的一部分。用户输入里即使包含引号、空格、OR,也只会被当成普通字符串。
注意:表名、字段名、排序字段不能直接用参数绑定。它们必须使用白名单。
allowed_order_fields = {"id", "created_at", "name"}
order_by = "created_at"
if order_by not in allowed_order_fields:
raise ValueError("排序字段非法")
sql = f"SELECT id, name FROM users ORDER BY {order_by} DESC"查询结果怎么读取
cursor.execute("SELECT id, name FROM users WHERE id = ?", (1,))
one_user = cursor.fetchone()
cursor.execute("SELECT id, name FROM users")
ten_users = cursor.fetchmany(10)
cursor.execute("SELECT id, name FROM users")
all_users = cursor.fetchall()| 方法 | 说明 | 风险 |
|---|---|---|
fetchone() | 读取一条 | 适合按主键查询 |
fetchmany(n) | 分批读取 | 适合较大结果集 |
fetchall() | 一次读取所有 | 大数据量可能撑爆内存 |
生产中不要对大表直接 fetchall()。如果需要导出几十万行,应该分页、游标分批或流式处理。
事务是什么
事务表示一组数据库操作要么全部成功,要么全部失败。典型例子是转账:
- A 账户扣 100。
- B 账户加 100。
- 写入转账流水。
这三步必须一起成功。如果只扣款没加款,数据就错了。
flowchart TD
A["开始事务"] --> B["扣减 A 账户"]
B --> C["增加 B 账户"]
C --> D["写入流水"]
D --> E{"三步是否都成功"}
E -- "成功" --> F["commit 提交"]
E -- "失败" --> G["rollback 回滚"]Python 示例:
try:
cursor.execute("UPDATE accounts SET balance = balance - ? WHERE id = ?", (100, 1))
cursor.execute("UPDATE accounts SET balance = balance + ? WHERE id = ?", (100, 2))
cursor.execute("INSERT INTO transfer_logs (from_id, to_id, amount) VALUES (?, ?, ?)", (1, 2, 100))
conn.commit()
except Exception:
conn.rollback()
raise如果忘记回滚,连接可能带着一个失败的事务继续被使用。连接池场景下,这会污染后续请求。
事务隔离级别
多个事务并发执行时,数据库要决定“一个事务能看到另一个事务的哪些变化”。这就是隔离级别。
| 现象 | 含义 | 例子 |
|---|---|---|
| 脏读 | 读到别人未提交的数据 | A 改余额未提交,B 已读到 |
| 不可重复读 | 同一事务两次读同一行结果不同 | B 第一次读 100,A 提交后 B 第二次读 80 |
| 幻读 | 同一范围查询两次行数不同 | 第一次查 10 条,别人插入后第二次查 11 条 |
常见隔离级别:
| 隔离级别 | 能避免什么 | 代价 |
|---|---|---|
| Read Uncommitted | 几乎不避免 | 数据一致性差 |
| Read Committed | 避免脏读 | 可能不可重复读 |
| Repeatable Read | 避免脏读、不可重复读 | 锁和版本管理更复杂 |
| Serializable | 最强隔离 | 并发性能最低 |
应用层要理解:隔离级别越高,不代表系统越好。商业项目要在一致性和并发性能之间取平衡。
上下文管理器封装事务
为了避免忘记提交、回滚、关闭连接,可以封装事务上下文。
import sqlite3
from contextlib import contextmanager
@contextmanager
def transaction(db_path: str):
conn = sqlite3.connect(db_path)
try:
yield conn
conn.commit()
except Exception:
conn.rollback()
raise
finally:
conn.close()
with transaction("demo.db") as conn:
conn.execute("INSERT INTO users (name, age) VALUES (?, ?)", ("李四", 20))
conn.execute("INSERT INTO users (name, age) VALUES (?, ?)", ("王五", 22))这样做的好处是业务代码不容易遗漏资源释放。
连接池为什么需要
数据库连接不是普通对象。创建连接通常涉及 TCP 连接、认证、权限校验、初始化会话等步骤。每个请求都新建连接会很慢。
连接池流程:
flowchart TD
A["应用启动"] --> B["创建连接池"]
C["请求到来"] --> D{"池中是否有空闲连接"}
D -- "有" --> E["借出连接"]
D -- "没有" --> F{"是否达到最大连接数"}
F -- "未达到" --> G["创建新连接"]
F -- "已达到" --> H["等待或超时"]
E --> I["执行 SQL"]
G --> I
I --> J["提交或回滚事务"]
J --> K["归还连接"]连接池不能无限大,原因:
- 数据库最大连接数有限。
- 每条连接都会占用数据库内存。
- 并发 SQL 太多会增加锁竞争、IO 压力和 CPU 调度。
- Web 服务多 worker 时,每个 worker 可能都有自己的连接池。
简单估算:
总连接数 = Web worker 数 * 每个进程连接池大小如果 4 个 worker,每个连接池最大 20,就是最多 80 条连接。还要给后台任务、管理工具、其他服务留余量。
连接池参数怎么理解
以 SQLAlchemy 为例,常见连接池配置如下:
from sqlalchemy import create_engine
engine = create_engine(
"mysql+pymysql://user:password@127.0.0.1:3306/app",
pool_size=10,
max_overflow=5,
pool_timeout=30,
pool_recycle=1800,
pool_pre_ping=True,
)| 参数 | 含义 | 配错后果 |
|---|---|---|
pool_size | 池里长期保留的连接数 | 太小会等待,太大会占用数据库连接 |
max_overflow | 池满后临时额外创建的连接数 | 太大可能瞬间打爆数据库 |
pool_timeout | 获取连接等待多久超时 | 太长会让请求长时间卡住 |
pool_recycle | 连接使用多久后回收重建 | 太长可能遇到数据库主动断开 |
pool_pre_ping | 借出连接前先探活 | 能减少“拿到死连接”问题 |
连接池状态变化:
flowchart TD
A["请求需要数据库"] --> B{"空闲连接是否可用"}
B -->|有| C["借出空闲连接"]
B -->|没有| D{"连接数是否小于 pool_size + max_overflow"}
D -->|是| E["创建临时连接"]
D -->|否| F["等待 pool_timeout"]
F --> G{"等待期间是否有连接归还"}
G -->|有| C
G -->|没有| H["抛出获取连接超时"]
C --> I["执行 SQL 和事务"]
E --> I
I --> J["提交或回滚"]
J --> K["归还或关闭连接"]生产估算不要只看应用自己:
数据库最大连接数
> Web 实例数 * worker 数 * 每个 worker 连接池上限
+ 定时任务连接
+ 后台消费者连接
+ 运维和监控连接
+ 预留余量例如 3 台机器,每台 4 个 worker,每个 worker pool_size=10,max_overflow=5,理论峰值就是:
3 * 4 * (10 + 5) = 180 条应用连接如果数据库最大连接数只有 200,再加上任务和运维连接,就很容易耗尽。
连接池耗尽时不要只盲目调大连接池。先排查:
- 是否有连接未关闭。
- 是否慢 SQL 占住连接。
- 是否事务过长。
- 是否外部接口调用放在事务里。
- 是否并发任务突然放大。
- 数据库本身是否 CPU、IO、锁等待严重。
SQLAlchemy Core、ORM、Session
SQLAlchemy 有两种常见用法:
| 方式 | 特点 | 适合 |
|---|---|---|
| Core | 更接近 SQL,用表达式构造查询 | 复杂 SQL、报表、批处理 |
| ORM | 用类映射表,用对象表达数据 | 业务系统、领域模型 |
ORM 示例:
from sqlalchemy import Integer, String, create_engine, select
from sqlalchemy.orm import DeclarativeBase, Mapped, Session, mapped_column
class Base(DeclarativeBase):
pass
class User(Base):
__tablename__ = "users"
id: Mapped[int] = mapped_column(Integer, primary_key=True)
name: Mapped[str] = mapped_column(String(50))
age: Mapped[int] = mapped_column(Integer)
engine = create_engine("sqlite:///demo.db", echo=True)
Base.metadata.create_all(engine)
with Session(engine) as session:
user = User(name="赵六", age=28)
session.add(user)
session.commit()
with Session(engine) as session:
users = session.scalars(select(User).where(User.age >= 18)).all()
for user in users:
print(user.id, user.name)Session 不是简单的数据库连接。它负责:
- 管理对象状态。
- 维护一次工作单元。
- 决定何时 flush SQL 到数据库。
- 控制事务提交和回滚。
- 从连接池借还连接。
不要把同一个 Session 当成全局变量复用。Web 项目通常是“每个请求一个 Session”。
SQLAlchemy Session 的状态流转
初学者很容易误解:session.add(user) 不是立刻把数据永久写入数据库,commit() 才是提交事务。中间还有 flush 这个动作。
flowchart TD
A["创建 Python 对象"] --> B["session.add"]
B --> C["对象进入 pending 状态"]
C --> D{"是否 flush"}
D -->|是| E["生成并发送 INSERT/UPDATE SQL"]
E --> F["数据库事务中已有修改"]
F --> G{"是否 commit"}
G -->|是| H["事务提交,其他事务可见"]
G -->|否,发生异常| I["rollback 回滚"]
D -->|否| J["等待查询或 commit 触发自动 flush"]
J --> E几个关键概念:
| 概念 | 含义 |
|---|---|
add | 把对象交给 Session 管理 |
flush | 把内存中的变更同步成 SQL 发给数据库,但事务还没提交 |
commit | 提交事务,让修改真正生效 |
rollback | 回滚事务,撤销未提交修改 |
close | 关闭 Session,归还连接和释放资源 |
示例:
with Session(engine) as session:
user = User(name="张三", age=18)
session.add(user)
session.flush()
print(user.id) # flush 后通常能拿到数据库生成的主键
session.commit()为什么有时还没 commit 就能拿到 id?因为 flush 已经执行了 INSERT,数据库生成了主键,但事务还没有提交。如果随后 rollback,这条数据仍然不会成为最终有效数据。
autoflush 是什么
SQLAlchemy 默认可能在查询前自动 flush,避免查询结果和 Session 内存状态不一致。
with Session(engine) as session:
user = User(name="张三", age=18)
session.add(user)
users = session.scalars(select(User)).all() # 查询前可能触发 autoflush这解释了一个现象:你以为只是查询,日志里却先出现了 INSERT。这不代表已经提交,只是 flush 到当前事务。
expire_on_commit 是什么
默认情况下,commit() 后对象可能过期。再次访问属性时,SQLAlchemy 可能重新查询数据库,以保证拿到最新值。
from sqlalchemy.orm import sessionmaker
SessionLocal = sessionmaker(bind=engine, expire_on_commit=False)Web 项目里有时会设置 expire_on_commit=False,避免提交后对象属性访问触发额外 SQL。但这也意味着你要理解对象可能不是数据库最新值。
SQLAlchemy 常见坑:N+1 查询
N+1 查询指:先查 1 次主表,再在循环里为每条记录查 1 次关联表。
错误思路:
orders = session.scalars(select(Order).limit(100)).all()
for order in orders:
print(order.user.name) # 每个 order 可能再查一次 user如果订单 100 条,可能执行 1 + 100 次 SQL。数据量一上来接口就会慢。
优化思路:
from sqlalchemy.orm import selectinload
orders = session.scalars(
select(Order).options(selectinload(Order.user)).limit(100)
).all()或者使用显式 JOIN,把需要的数据一次查出来。N+1 的本质不是 ORM 的错,而是你没有控制关联数据加载方式。
FastAPI 中的数据库 Session
典型写法:
from typing import Generator
from fastapi import Depends, FastAPI
from sqlalchemy import create_engine, select
from sqlalchemy.orm import Session, sessionmaker
engine = create_engine("sqlite:///demo.db")
SessionLocal = sessionmaker(bind=engine, autoflush=False, autocommit=False)
app = FastAPI()
def get_db() -> Generator[Session, None, None]:
db = SessionLocal()
try:
yield db
finally:
db.close()
@app.get("/users")
def list_users(db: Session = Depends(get_db)):
return db.execute(select(User)).scalars().all()为什么用 yield?因为请求结束后 FastAPI 会继续执行 finally,确保 Session 被关闭,连接归还连接池。
FastAPI 请求里的事务边界
很多项目会在依赖函数里只负责创建和关闭 Session,把 commit/rollback 放在 Service 层。这样做更清晰,因为 Service 才知道一个完整业务动作到哪里结束。
推荐流程:
flowchart TD
A["HTTP 请求进入 API"] --> B["Depends 创建 Session"]
B --> C["API 调用 Service"]
C --> D["Service 执行业务校验"]
D --> E["Repository 执行多条 SQL"]
E --> F{"业务是否全部成功"}
F -->|是| G["Service commit"]
F -->|否| H["Service rollback"]
G --> I["API 返回成功响应"]
H --> J["异常处理返回错误响应"]
I --> K["Depends finally 关闭 Session"]
J --> K示例:
def create_asset_service(db: Session, request: AssetCreate) -> Asset:
try:
asset = asset_repository.insert_asset(db, request)
audit_repository.insert_log(db, "CREATE_ASSET", asset.id)
db.commit()
db.refresh(asset)
return asset
except Exception:
db.rollback()
raise为什么不建议 Repository 自己 commit?
- Repository 只知道“单次数据访问”,不知道完整业务是否结束。
- 一个业务动作可能要写多张表。
- 如果每个 Repository 都 commit,中途失败时前面的修改无法一起回滚。
- 事务边界分散后,排查数据不一致非常痛苦。
外部接口不要放在长事务里
错误示例:
def create_asset(db: Session, request: dict):
asset_repository.insert_asset(db, request)
ai_client.analyze_asset(request) # 外部接口可能耗时 5 秒
audit_repository.insert_log(db, "CREATE_ASSET")
db.commit()问题:数据库事务在等待 AI 接口时一直开着,可能持有锁和连接。并发上来后会出现锁等待、连接池耗尽、接口整体变慢。
更好的方案:
def create_asset(db: Session, request: dict):
try:
asset = asset_repository.insert_asset(db, request)
audit_repository.insert_log(db, "CREATE_ASSET", asset.id)
task_repository.insert_task(db, "ANALYZE_ASSET", asset.id)
db.commit()
return asset
except Exception:
db.rollback()
raise事务里只写本地数据库关键状态。AI 分析、消息发送、第三方通知可以通过任务表、消息队列或后台任务异步处理。
Repository 和 Service 的事务边界
推荐原则:
- Repository 只负责单表或少量数据访问。
- Service 负责一个完整业务动作。
- 事务边界通常放在 Service 层。
错误做法:
def create_order():
order_repository.save_order()
order_repository.commit()
stock_repository.deduct_stock()
stock_repository.commit()如果订单提交成功后库存扣减失败,数据就不一致。
正确思路:
def create_order(db: Session, user_id: int, sku_id: int):
try:
order_repository.save_order(db, user_id, sku_id)
stock_repository.deduct_stock(db, sku_id)
audit_repository.save_log(db, "CREATE_ORDER")
db.commit()
except Exception:
db.rollback()
raise一个业务动作中的多次数据库修改应该在同一个事务里提交。
Alembic 迁移
表结构会变,例如新增字段、索引、表。不要靠手工改线上数据库,应该用迁移文件记录版本。
常见流程:
flowchart TD
A["修改 ORM 模型或手写迁移意图"] --> B["生成迁移文件"]
B --> C["人工检查 SQL 是否正确"]
C --> D["测试环境执行升级"]
D --> E["生产发布执行迁移"]
E --> F["必要时按回滚脚本降级"]命令示例:
alembic init migrations
alembic revision -m "add asset table"
alembic upgrade head
alembic downgrade -1迁移文件要进 Git。否则代码和数据库结构会失去对应关系,部署时容易出现“代码需要字段,但线上表没有字段”的问题。
迁移为什么必须人工检查
自动生成迁移很方便,但不能无脑执行。因为工具只能比较“模型结构变化”,并不真正理解业务数据和线上风险。
常见风险:
| 变更 | 风险 | 更稳做法 |
|---|---|---|
| 新增非空字段 | 老数据没有值,迁移失败 | 先加可空字段,回填数据,再改非空 |
| 删除字段 | 代码或报表还在用 | 先停止使用,观察后再删除 |
| 修改字段类型 | 可能转换失败或锁表 | 评估数据量,分批迁移 |
| 新增大表索引 | 可能长时间锁表或占 IO | 低峰执行,使用在线 DDL 能力 |
| 重命名字段 | 部署期间新旧代码不兼容 | 兼容式发布,双写或双读过渡 |
兼容式字段变更流程:
flowchart TD
A["第 1 次发布:新增可空字段"] --> B["代码同时兼容新旧字段"]
B --> C["后台脚本回填历史数据"]
C --> D["校验新字段数据完整"]
D --> E["第 2 次发布:代码只读新字段"]
E --> F["第 3 次迁移:删除旧字段或加非空约束"]这比“一次迁移改完”麻烦,但生产风险小很多。数据库变更和代码发布要当作一个整体设计,不要只看本地迁移能不能跑。
分页、大查询和批处理
普通分页:
SELECT id, name FROM assets ORDER BY id LIMIT 20 OFFSET 40;大偏移量分页可能越来越慢,因为数据库需要跳过很多行。大表推荐基于游标的分页:
SELECT id, name
FROM assets
WHERE id > ?
ORDER BY id
LIMIT 100;Python 分批处理:
last_id = 0
while True:
cursor.execute(
"SELECT id, name FROM assets WHERE id > ? ORDER BY id LIMIT 100",
(last_id,),
)
rows = cursor.fetchall()
if not rows:
break
for row in rows:
print(row)
last_id = rows[-1][0]不要一次性把几百万行加载进内存再处理。数据量越大,越要分批、断点续跑、记录进度。
慢 SQL 排查
慢 SQL 排查流程:
flowchart TD
A["发现接口或任务慢"] --> B["定位具体 SQL"]
B --> C["查看执行耗时和扫描行数"]
C --> D["执行 EXPLAIN"]
D --> E{"是否走合适索引"}
E -- "否" --> F["补索引或改写条件"]
E -- "是" --> G{"返回行数是否太多"}
G -- "是" --> H["分页 限字段 增加筛选条件"]
G -- "否" --> I{"是否锁等待"}
I -- "是" --> J["查事务和锁持有者"]
I -- "否" --> K["检查排序 分组 JOIN 临时表"]排查时要收集:
| 信息 | 作用 |
|---|---|
| SQL 文本 | 确认真正执行的语句 |
| 参数 | 同一 SQL 不同参数性能可能不同 |
| 执行时间 | 判断慢的严重程度 |
| 扫描行数 | 判断是否大量无效扫描 |
| 执行计划 | 判断索引、Join、排序方式 |
| 锁等待 | 判断是否被其他事务阻塞 |
不要一慢就加缓存。缓存只能绕过部分读请求,不能解决错误索引、事务锁、SQL 写法和数据模型问题。
常见生产问题
| 问题 | 原因 | 后果 | 解决 |
|---|---|---|---|
| 连接耗尽 | 连接未关闭、池太小、慢 SQL 占住连接 | 新请求获取不到连接 | 关闭连接、调池大小、优化慢 SQL |
| 事务太长 | 事务里调用外部接口或处理大文件 | 锁持有时间长 | 事务只包数据库关键操作 |
fetchall() 内存高 | 一次读太多数据 | OOM 或进程被杀 | 分页、分批、流式 |
| SQL 注入 | 拼接用户输入 | 数据泄露、被删表 | 参数绑定、白名单 |
| N+1 查询 | 循环中逐条查关联数据 | 接口慢、数据库压力大 | JOIN、批量查询、预加载 |
| 迁移丢失 | 手工改库没记录 | 环境不一致 | Alembic 管理迁移 |
| Session 全局复用 | 不理解 Session 生命周期 | 脏数据、事务串扰 | 每请求一个 Session |
商业场景:资产导入和审计日志
医疗数据资产平台导入资产时,一次请求可能要做:
- 校验来源系统是否合法。
- 写入资产主表。
- 写入字段明细表。
- 写入审计日志。
- 更新导入批次状态。
这些数据库修改应该在一个事务里。如果第 3 步失败,前面的资产主表也要回滚,否则会出现“资产存在但字段不完整”的脏数据。
简化代码:
def import_asset(db: Session, request: dict):
try:
asset_id = asset_repository.insert_asset(db, request)
field_repository.batch_insert_fields(db, asset_id, request["fields"])
audit_repository.insert_log(db, "IMPORT_ASSET", asset_id)
batch_repository.mark_success(db, request["batch_id"])
db.commit()
return asset_id
except Exception:
db.rollback()
raise如果字段很多,可以分批插入,但事务边界要根据业务一致性要求决定。强一致要求高时同事务;数据量极大时可能要改成“批次状态 + 分段提交 + 失败补偿”。
商业 Demo:资产导入批次的 Repository 和事务
下面用 SQLAlchemy 演示一个更完整的导入批次流程。目标不是堆代码,而是看清楚商业项目里数据库访问怎么分层。
业务规则:
- 每次导入先创建一个批次。
- 每条资产按
asset_code去重。 - 成功导入资产后写审计日志。
- 单条资产失败时记录失败原因。
- 批次状态最后更新为
SUCCESS或PARTIAL_FAILED。 - 每一批数据在一个事务中提交,失败则整批回滚。
简化模型:
from sqlalchemy import ForeignKey, Integer, String, create_engine, select
from sqlalchemy.orm import DeclarativeBase, Mapped, Session, mapped_column, sessionmaker
class Base(DeclarativeBase):
pass
class ImportBatch(Base):
__tablename__ = "import_batches"
id: Mapped[int] = mapped_column(Integer, primary_key=True)
file_name: Mapped[str] = mapped_column(String(200))
status: Mapped[str] = mapped_column(String(30))
class Asset(Base):
__tablename__ = "assets"
id: Mapped[int] = mapped_column(Integer, primary_key=True)
asset_code: Mapped[str] = mapped_column(String(50), unique=True)
asset_name: Mapped[str] = mapped_column(String(100))
department: Mapped[str] = mapped_column(String(100))
class ImportErrorRow(Base):
__tablename__ = "import_error_rows"
id: Mapped[int] = mapped_column(Integer, primary_key=True)
batch_id: Mapped[int] = mapped_column(ForeignKey("import_batches.id"))
line_no: Mapped[int] = mapped_column(Integer)
reason: Mapped[str] = mapped_column(String(300))
class AuditLog(Base):
__tablename__ = "audit_logs"
id: Mapped[int] = mapped_column(Integer, primary_key=True)
action: Mapped[str] = mapped_column(String(50))
target_id: Mapped[int] = mapped_column(Integer)Repository:
class BatchRepository:
def create(self, db: Session, file_name: str) -> ImportBatch:
batch = ImportBatch(file_name=file_name, status="RUNNING")
db.add(batch)
db.flush()
return batch
def mark_finished(self, db: Session, batch: ImportBatch, has_error: bool) -> None:
batch.status = "PARTIAL_FAILED" if has_error else "SUCCESS"
class AssetRepository:
def exists_by_code(self, db: Session, asset_code: str) -> bool:
stmt = select(Asset.id).where(Asset.asset_code == asset_code)
return db.scalar(stmt) is not None
def create(self, db: Session, row: dict) -> Asset:
asset = Asset(
asset_code=row["asset_code"],
asset_name=row["asset_name"],
department=row["department"],
)
db.add(asset)
db.flush()
return asset
class ErrorRowRepository:
def create(self, db: Session, batch_id: int, line_no: int, reason: str) -> None:
db.add(ImportErrorRow(batch_id=batch_id, line_no=line_no, reason=reason))
class AuditRepository:
def create(self, db: Session, action: str, target_id: int) -> None:
db.add(AuditLog(action=action, target_id=target_id))Service:
class AssetImportService:
def __init__(
self,
batch_repo: BatchRepository,
asset_repo: AssetRepository,
error_repo: ErrorRowRepository,
audit_repo: AuditRepository,
):
self.batch_repo = batch_repo
self.asset_repo = asset_repo
self.error_repo = error_repo
self.audit_repo = audit_repo
def import_rows(self, db: Session, file_name: str, rows: list[dict]) -> int:
try:
batch = self.batch_repo.create(db, file_name)
has_error = False
for line_no, row in enumerate(rows, start=2):
try:
self._validate(row)
if self.asset_repo.exists_by_code(db, row["asset_code"]):
raise ValueError(f"资产编号重复: {row['asset_code']}")
asset = self.asset_repo.create(db, row)
self.audit_repo.create(db, "IMPORT_ASSET", asset.id)
except ValueError as exc:
has_error = True
self.error_repo.create(db, batch.id, line_no, str(exc))
self.batch_repo.mark_finished(db, batch, has_error)
db.commit()
return batch.id
except Exception:
db.rollback()
raise
def _validate(self, row: dict) -> None:
if not row.get("asset_code"):
raise ValueError("asset_code 不能为空")
if not row.get("asset_name"):
raise ValueError("asset_name 不能为空")
if not row.get("department"):
raise ValueError("department 不能为空")运行 Demo:
engine = create_engine("sqlite:///asset_import.db", echo=True)
Base.metadata.create_all(engine)
SessionLocal = sessionmaker(bind=engine, autoflush=False)
service = AssetImportService(
BatchRepository(),
AssetRepository(),
ErrorRowRepository(),
AuditRepository(),
)
rows = [
{"asset_code": "A001", "asset_name": "心电监护仪", "department": "ICU"},
{"asset_code": "", "asset_name": "输液泵", "department": "急诊"},
]
with SessionLocal() as db:
batch_id = service.import_rows(db, "assets.csv", rows)
print("batch_id=", batch_id)流程图:
flowchart TD
A["Service 开始导入"] --> B["创建 ImportBatch"]
B --> C["遍历 CSV 行"]
C --> D{"字段是否合法"}
D -->|否| E["写 ImportErrorRow"]
D -->|是| F{"asset_code 是否已存在"}
F -->|是| E
F -->|否| G["写 Asset"]
G --> H["写 AuditLog"]
E --> I{"是否还有下一行"}
H --> I
I -->|有| C
I -->|无| J["更新批次状态"]
J --> K{"整体是否异常"}
K -->|否| L["commit"]
K -->|是| M["rollback"]这个 Demo 要看懂几个关键点:
- Repository 不提交事务,只操作数据库对象。
- Service 负责事务边界,因为它知道一个完整业务动作包含哪些步骤。
- 单行数据校验失败被记录到错误表,不一定导致整个批次失败。
- 系统异常会回滚整个事务,例如数据库断开、表结构错误。
db.flush()用来提前拿到主键,方便写关联表,但不是提交事务。
如果数据量非常大,不建议几十万行放在一个事务里。可以按批次分段提交,并使用批次状态、错误表和幂等键保证可恢复。
数据库访问的可观测性
线上排查数据库问题时,只知道“接口慢”不够。至少要能看到:
| 指标或日志 | 作用 |
|---|---|
| SQL 耗时 | 判断慢 SQL |
| 连接池借出等待时间 | 判断连接池是否紧张 |
| 事务持续时间 | 判断是否长事务 |
| commit/rollback 次数 | 判断失败比例 |
| 影响行数 | 判断更新是否命中 |
| 锁等待或死锁错误 | 判断并发冲突 |
简单 SQL 计时包装:
import time
import logging
logger = logging.getLogger(__name__)
def execute_with_log(cursor, sql: str, params: tuple = ()):
start = time.perf_counter()
try:
cursor.execute(sql, params)
return cursor
finally:
cost_ms = int((time.perf_counter() - start) * 1000)
logger.info("sql cost_ms=%s rowcount=%s sql=%s", cost_ms, cursor.rowcount, sql)生产日志不要直接打印敏感参数,例如手机号、身份证号、Token、密码。可以记录 SQL 模板、耗时、业务 ID、request_id 和影响行数。
面试标准回答
Python 操作数据库的基本流程是什么?
先通过数据库驱动创建连接,再创建 Cursor 或 Session,执行带参数的 SQL,读取结果。如果是写操作,要在业务成功时提交事务,失败时回滚事务,最后关闭或归还连接。
为什么不能拼接 SQL?
因为用户输入如果被当成 SQL 语法的一部分,就可能改变原 SQL 逻辑,造成 SQL 注入。参数绑定会把 SQL 模板和参数分开处理,用户输入只作为值,不作为语法执行。字段名和排序字段不能参数绑定,要用白名单。
连接池为什么不能无限大?
连接池能复用连接、减少建连成本,但连接会占用数据库资源。连接过多会打满数据库最大连接数,增加锁竞争和 CPU 调度。多 worker 部署时要按“worker 数乘以池大小”估算总连接数。
SQLAlchemy Session 是什么?
Session 是一次工作单元,不只是连接。它管理对象状态、事务、flush、commit、rollback,并从连接池获取连接。Web 项目通常每个请求创建一个 Session,请求结束后关闭,不能全局复用。
慢 SQL 怎么排查?
先定位具体 SQL 和参数,再看耗时、扫描行数、执行计划、索引使用、排序分组、JOIN 和锁等待。不要先加缓存,要先判断是索引问题、返回数据过多、事务阻塞还是 SQL 写法问题。
flush 和 commit 有什么区别?flush 是把 Session 中的对象变更转换成 SQL 发给数据库,数据库在当前事务里已经能生成主键和约束检查,但事务还没有提交;commit 是提交事务,让修改真正生效并对其他事务可见。如果 flush 后 rollback,数据仍然会被回滚。
Repository 为什么不建议自己 commit?
Repository 只负责数据访问,不知道完整业务动作是否结束。一个业务动作可能要写订单、库存、日志多张表,如果每个 Repository 自己 commit,中途失败就无法整体回滚。事务边界通常应放在 Service 层。
为什么外部接口不要放在数据库事务里?
外部接口可能慢、超时或不稳定。如果事务里等待外部接口,会长时间占用数据库连接和锁,导致连接池耗尽、锁等待和接口变慢。更好的做法是在本地事务里写任务表或消息,再异步调用外部系统。
数据库连接池耗尽怎么排查?
先看是否有连接未关闭,再看慢 SQL、长事务、锁等待、连接池配置、worker 数量和并发任务。不要第一反应就把连接池调大,因为调大可能把压力转移到数据库,导致数据库最大连接数耗尽。
数据库迁移为什么不能直接线上手工改表?
手工改表没有版本记录,容易造成代码和数据库结构不一致。迁移文件能让开发、测试、生产环境保持同一演进历史,也便于回滚和审计。复杂变更还要分阶段兼容发布。
关联知识点
小结
Python 数据库访问的主线是连接、SQL、结果、事务、连接池和 ORM。零基础阶段要先理解驱动怎么执行 SQL,再理解事务为什么能保证一致性,最后再学 SQLAlchemy 和迁移工具。商业项目中最重要的不是“能查出来”,而是能保证数据一致、性能可控、资源可释放、问题可排查。
