Skip to content

Python 数据库访问

大多数后端程序最终都要保存数据,例如用户、订单、采集任务、审计日志、AI 对话记录、资产目录。Python 访问数据库不要只会 execute(),还要理解连接、游标、事务、连接池、ORM、迁移和慢 SQL 排查。

本页以关系型数据库为主,例如 MySQL、PostgreSQL、SQLite。不同数据库语法细节不同,但 Python 访问数据库的核心流程相通。

学习目标

学完本页你应该能回答:

  1. Python 是怎么通过驱动访问数据库的。
  2. Connection、Cursor、Transaction 分别是什么。
  3. 为什么参数绑定能防 SQL 注入。
  4. 事务什么时候提交、什么时候回滚。
  5. 连接池为什么需要,为什么不能无限大。
  6. SQLAlchemy Core、ORM、Session、Alembic 分别解决什么问题。
  7. 大查询、慢 SQL、锁等待、连接耗尽怎么排查。
  8. 商业项目中 Repository 和 Service 的事务边界怎么划分。

数据库访问的基本流程

mermaid
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 发给数据库服务,由数据库完成解析、优化、执行、锁控制和事务处理。

常见方案怎么选

方案说明适合场景
sqlite3Python 标准库自带,单文件数据库学习、小工具、本地缓存
PyMySQL / mysqlclientMySQL 驱动连接 MySQL
psycopgPostgreSQL 驱动连接 PostgreSQL
SQLAlchemy Core用 Python API 构造 SQL想减少手写 SQL,又要控制 SQL
SQLAlchemy ORM类和表映射,通过对象操作数据中大型 Web 项目
Alembic数据库迁移工具管理表结构版本

学习路线:

  1. 先用 sqlite3 理解连接、游标、SQL、事务。
  2. 再用 MySQL/PostgreSQL 驱动理解网络数据库。
  3. 再学 SQLAlchemy Session 生命周期和事务边界。
  4. 最后学 Alembic,把表结构变更纳入版本管理。

DB-API:Connection 和 Cursor

Python 数据库驱动大多遵循 DB-API 思路。

对象作用类比
Connection到数据库的一条连接,承载事务一条电话线
Cursor执行 SQL、读取结果的游标接线员
Transaction一组要么全成功要么全失败的操作一份完整业务单据

最小 Demo:

python
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()

这段代码的执行过程:

  1. sqlite3.connect() 打开数据库连接。
  2. conn.cursor() 创建游标。
  3. cursor.execute() 把 SQL 发送给数据库。
  4. 插入、更新、删除后要 commit(),否则事务可能没有真正提交。
  5. 查询通过 fetchone()fetchmany()fetchall() 获取结果。
  6. 最后关闭游标和连接。

驱动到底做了什么

Python 代码不会直接“进入数据库表里拿数据”。真正的链路通常是:Python 调用驱动,驱动把 SQL 和参数按照数据库协议发给数据库服务,数据库执行完成后再把结果返回给驱动。

mermaid
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 对象"]

这条链路解释了几个常见现象:

  1. SQL 写错不是 Python 解释器发现的,而是数据库返回语法错误。
  2. 网络数据库连接失败可能是账号、密码、端口、防火墙、网络、TLS、权限任意一层出问题。
  3. 查询慢不一定是 Python 慢,通常要看数据库执行计划、锁等待、IO、索引和返回行数。
  4. 驱动会把数据库类型转换成 Python 类型,例如 SQL 的 VARCHAR 变成 strINT 变成 int,时间类型变成 datetime

SQLite 比较特殊,它是进程内文件数据库,不需要 TCP 连接数据库服务;MySQL、PostgreSQL 则通常是通过网络连接数据库进程。

数据库连接方式适合
SQLite本地文件学习、本地工具、轻量缓存
MySQLTCP 连接数据库服务Web 业务、订单、资产、权限
PostgreSQLTCP 连接数据库服务复杂查询、GIS、强 SQL 能力

execute、executemany 和影响行数

单条执行:

python
cursor.execute(
    "INSERT INTO users (name, age) VALUES (?, ?)",
    ("张三", 18),
)

批量执行:

python
rows = [
    ("张三", 18),
    ("李四", 20),
    ("王五", 22),
]
cursor.executemany("INSERT INTO users (name, age) VALUES (?, ?)", rows)

executemany 的价值是减少 Python 和数据库之间的来回通信次数。但它也不是“越大越好”,一次批量太大可能导致:

  1. SQL 包太大。
  2. 事务持有锁太久。
  3. 失败回滚成本高。
  4. 内存占用升高。

商业批处理建议按批次提交,例如每 500 或 1000 行一批,具体要结合数据库性能、锁竞争和失败补偿策略。

影响行数:

python
cursor.execute("UPDATE users SET age = ? WHERE id = ?", (19, 1))
print(cursor.rowcount)

rowcount 常用于判断更新是否真正命中数据。例如订单支付更新状态时,如果影响行数为 0,可能说明订单不存在、状态已经变化或并发抢占失败。

参数绑定为什么能防 SQL 注入

危险写法:

python
name = "' OR '1'='1"
sql = f"SELECT * FROM users WHERE name = '{name}'"
cursor.execute(sql)

最终 SQL 可能变成:

sql
SELECT * FROM users WHERE name = '' OR '1'='1'

这会绕过原本的过滤条件。

正确写法:

python
cursor.execute("SELECT * FROM users WHERE name = ?", (name,))

MySQL 常见参数符号:

python
cursor.execute("SELECT * FROM users WHERE name = %s", (name,))

参数绑定的原理是:SQL 模板和参数分开发给驱动,驱动把用户输入当作“值”处理,而不是当作 SQL 语法的一部分。用户输入里即使包含引号、空格、OR,也只会被当成普通字符串。

注意:表名、字段名、排序字段不能直接用参数绑定。它们必须使用白名单。

python
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"

查询结果怎么读取

python
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()。如果需要导出几十万行,应该分页、游标分批或流式处理。

事务是什么

事务表示一组数据库操作要么全部成功,要么全部失败。典型例子是转账:

  1. A 账户扣 100。
  2. B 账户加 100。
  3. 写入转账流水。

这三步必须一起成功。如果只扣款没加款,数据就错了。

mermaid
flowchart TD
    A["开始事务"] --> B["扣减 A 账户"]
    B --> C["增加 B 账户"]
    C --> D["写入流水"]
    D --> E{"三步是否都成功"}
    E -- "成功" --> F["commit 提交"]
    E -- "失败" --> G["rollback 回滚"]

Python 示例:

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最强隔离并发性能最低

应用层要理解:隔离级别越高,不代表系统越好。商业项目要在一致性和并发性能之间取平衡。

上下文管理器封装事务

为了避免忘记提交、回滚、关闭连接,可以封装事务上下文。

python
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 连接、认证、权限校验、初始化会话等步骤。每个请求都新建连接会很慢。

连接池流程:

mermaid
flowchart TD
    A["应用启动"] --> B["创建连接池"]
    C["请求到来"] --> D{"池中是否有空闲连接"}
    D -- "有" --> E["借出连接"]
    D -- "没有" --> F{"是否达到最大连接数"}
    F -- "未达到" --> G["创建新连接"]
    F -- "已达到" --> H["等待或超时"]
    E --> I["执行 SQL"]
    G --> I
    I --> J["提交或回滚事务"]
    J --> K["归还连接"]

连接池不能无限大,原因:

  1. 数据库最大连接数有限。
  2. 每条连接都会占用数据库内存。
  3. 并发 SQL 太多会增加锁竞争、IO 压力和 CPU 调度。
  4. Web 服务多 worker 时,每个 worker 可能都有自己的连接池。

简单估算:

text
总连接数 = Web worker 数 * 每个进程连接池大小

如果 4 个 worker,每个连接池最大 20,就是最多 80 条连接。还要给后台任务、管理工具、其他服务留余量。

连接池参数怎么理解

以 SQLAlchemy 为例,常见连接池配置如下:

python
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借出连接前先探活能减少“拿到死连接”问题

连接池状态变化:

mermaid
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["归还或关闭连接"]

生产估算不要只看应用自己:

text
数据库最大连接数
  > Web 实例数 * worker 数 * 每个 worker 连接池上限
  + 定时任务连接
  + 后台消费者连接
  + 运维和监控连接
  + 预留余量

例如 3 台机器,每台 4 个 worker,每个 worker pool_size=10,max_overflow=5,理论峰值就是:

text
3 * 4 * (10 + 5) = 180 条应用连接

如果数据库最大连接数只有 200,再加上任务和运维连接,就很容易耗尽。

连接池耗尽时不要只盲目调大连接池。先排查:

  1. 是否有连接未关闭。
  2. 是否慢 SQL 占住连接。
  3. 是否事务过长。
  4. 是否外部接口调用放在事务里。
  5. 是否并发任务突然放大。
  6. 数据库本身是否 CPU、IO、锁等待严重。

SQLAlchemy Core、ORM、Session

SQLAlchemy 有两种常见用法:

方式特点适合
Core更接近 SQL,用表达式构造查询复杂 SQL、报表、批处理
ORM用类映射表,用对象表达数据业务系统、领域模型

ORM 示例:

python
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 不是简单的数据库连接。它负责:

  1. 管理对象状态。
  2. 维护一次工作单元。
  3. 决定何时 flush SQL 到数据库。
  4. 控制事务提交和回滚。
  5. 从连接池借还连接。

不要把同一个 Session 当成全局变量复用。Web 项目通常是“每个请求一个 Session”。

SQLAlchemy Session 的状态流转

初学者很容易误解:session.add(user) 不是立刻把数据永久写入数据库,commit() 才是提交事务。中间还有 flush 这个动作。

mermaid
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,归还连接和释放资源

示例:

python
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 内存状态不一致。

python
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 可能重新查询数据库,以保证拿到最新值。

python
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 次关联表。

错误思路:

python
orders = session.scalars(select(Order).limit(100)).all()
for order in orders:
    print(order.user.name)  # 每个 order 可能再查一次 user

如果订单 100 条,可能执行 1 + 100 次 SQL。数据量一上来接口就会慢。

优化思路:

python
from sqlalchemy.orm import selectinload

orders = session.scalars(
    select(Order).options(selectinload(Order.user)).limit(100)
).all()

或者使用显式 JOIN,把需要的数据一次查出来。N+1 的本质不是 ORM 的错,而是你没有控制关联数据加载方式。

FastAPI 中的数据库 Session

典型写法:

python
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 才知道一个完整业务动作到哪里结束。

推荐流程:

mermaid
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

示例:

python
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?

  1. Repository 只知道“单次数据访问”,不知道完整业务是否结束。
  2. 一个业务动作可能要写多张表。
  3. 如果每个 Repository 都 commit,中途失败时前面的修改无法一起回滚。
  4. 事务边界分散后,排查数据不一致非常痛苦。

外部接口不要放在长事务里

错误示例:

python
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 接口时一直开着,可能持有锁和连接。并发上来后会出现锁等待、连接池耗尽、接口整体变慢。

更好的方案:

python
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 的事务边界

推荐原则:

  1. Repository 只负责单表或少量数据访问。
  2. Service 负责一个完整业务动作。
  3. 事务边界通常放在 Service 层。

错误做法:

python
def create_order():
    order_repository.save_order()
    order_repository.commit()
    stock_repository.deduct_stock()
    stock_repository.commit()

如果订单提交成功后库存扣减失败,数据就不一致。

正确思路:

python
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 迁移

表结构会变,例如新增字段、索引、表。不要靠手工改线上数据库,应该用迁移文件记录版本。

常见流程:

mermaid
flowchart TD
    A["修改 ORM 模型或手写迁移意图"] --> B["生成迁移文件"]
    B --> C["人工检查 SQL 是否正确"]
    C --> D["测试环境执行升级"]
    D --> E["生产发布执行迁移"]
    E --> F["必要时按回滚脚本降级"]

命令示例:

bash
alembic init migrations
alembic revision -m "add asset table"
alembic upgrade head
alembic downgrade -1

迁移文件要进 Git。否则代码和数据库结构会失去对应关系,部署时容易出现“代码需要字段,但线上表没有字段”的问题。

迁移为什么必须人工检查

自动生成迁移很方便,但不能无脑执行。因为工具只能比较“模型结构变化”,并不真正理解业务数据和线上风险。

常见风险:

变更风险更稳做法
新增非空字段老数据没有值,迁移失败先加可空字段,回填数据,再改非空
删除字段代码或报表还在用先停止使用,观察后再删除
修改字段类型可能转换失败或锁表评估数据量,分批迁移
新增大表索引可能长时间锁表或占 IO低峰执行,使用在线 DDL 能力
重命名字段部署期间新旧代码不兼容兼容式发布,双写或双读过渡

兼容式字段变更流程:

mermaid
flowchart TD
    A["第 1 次发布:新增可空字段"] --> B["代码同时兼容新旧字段"]
    B --> C["后台脚本回填历史数据"]
    C --> D["校验新字段数据完整"]
    D --> E["第 2 次发布:代码只读新字段"]
    E --> F["第 3 次迁移:删除旧字段或加非空约束"]

这比“一次迁移改完”麻烦,但生产风险小很多。数据库变更和代码发布要当作一个整体设计,不要只看本地迁移能不能跑。

分页、大查询和批处理

普通分页:

sql
SELECT id, name FROM assets ORDER BY id LIMIT 20 OFFSET 40;

大偏移量分页可能越来越慢,因为数据库需要跳过很多行。大表推荐基于游标的分页:

sql
SELECT id, name
FROM assets
WHERE id > ?
ORDER BY id
LIMIT 100;

Python 分批处理:

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 排查流程:

mermaid
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

商业场景:资产导入和审计日志

医疗数据资产平台导入资产时,一次请求可能要做:

  1. 校验来源系统是否合法。
  2. 写入资产主表。
  3. 写入字段明细表。
  4. 写入审计日志。
  5. 更新导入批次状态。

这些数据库修改应该在一个事务里。如果第 3 步失败,前面的资产主表也要回滚,否则会出现“资产存在但字段不完整”的脏数据。

简化代码:

python
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 演示一个更完整的导入批次流程。目标不是堆代码,而是看清楚商业项目里数据库访问怎么分层。

业务规则:

  1. 每次导入先创建一个批次。
  2. 每条资产按 asset_code 去重。
  3. 成功导入资产后写审计日志。
  4. 单条资产失败时记录失败原因。
  5. 批次状态最后更新为 SUCCESSPARTIAL_FAILED
  6. 每一批数据在一个事务中提交,失败则整批回滚。

简化模型:

python
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:

python
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:

python
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:

python
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)

流程图:

mermaid
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 要看懂几个关键点:

  1. Repository 不提交事务,只操作数据库对象。
  2. Service 负责事务边界,因为它知道一个完整业务动作包含哪些步骤。
  3. 单行数据校验失败被记录到错误表,不一定导致整个批次失败。
  4. 系统异常会回滚整个事务,例如数据库断开、表结构错误。
  5. db.flush() 用来提前拿到主键,方便写关联表,但不是提交事务。

如果数据量非常大,不建议几十万行放在一个事务里。可以按批次分段提交,并使用批次状态、错误表和幂等键保证可恢复。

数据库访问的可观测性

线上排查数据库问题时,只知道“接口慢”不够。至少要能看到:

指标或日志作用
SQL 耗时判断慢 SQL
连接池借出等待时间判断连接池是否紧张
事务持续时间判断是否长事务
commit/rollback 次数判断失败比例
影响行数判断更新是否命中
锁等待或死锁错误判断并发冲突

简单 SQL 计时包装:

python
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 和迁移工具。商业项目中最重要的不是“能查出来”,而是能保证数据一致、性能可控、资源可释放、问题可排查。