第35章:SQLAlchemy极端性能——百万行写入与只读分析负载
一、项目背景
“每天早上 6 点定时导入 50 万条商品数据——从 30 分钟优化到 3 分钟,还是不够快。客户要求 1 分钟内完成。”
星云电商的核心指标——每天凌晨的"第三方商品价格同步"作业——随着供应商数量从 100 家增长到 5000 家,需要处理的数据量从 50 万行增加到 300 万行。原有的 ORM 批量导入方案在 50 万行时可以 3 分钟跑完,但在 300 万行时膨胀到 18 分钟——远远超过了客户的 5 分钟 SLA。
性能瓶颈主要集中三个方面:
- 写入路径:即使使用了
session.execute(insert()),300 万条数据的网络往返和数据量(约 200MB)仍然是瓶颈。 - 连接池争抢:导入作业使用 10 个 Worker 并行写入——共享同一个连接池(pool_size=5),连接争抢导致排队。
- 只读分析查询被写入阻塞:同一 PostgreSQL 实例上有 OLTP 写入(导入作业)和 OLAP 查询(报表统计)混合运行——前者占用大量 I/O 和锁,导致后者延迟飙升。
优化方向涉及三个层面:
- 写入优化:COPY 命令(PostgreSQL)、多值 INSERT、无 ORM 追踪批量——逐级榨取数据库吞吐极限。
- 读写分离:将写入路由到主库,报表查询路由到只读副本——Session 绑定多个 Engine。
- 会话隔离:OLTP Session 关闭 expire、关闭 autoflush、使用合适的事务隔离级别;OLAP Session 使用单独的连接池。
本章将设计一个日增百万订单明细的导入作业,对比四种写入方式的吞吐,实现读写分离路由,并给出生产级容量规划建议。
二、项目设计
场景:凌晨 5:45,数据导入作业还有 30% 未完成,监控面板上 ETA 已经超过了 SLA 死线。大师带着小胖和小白在白板前分析火焰图。
小胖:“太离谱了——每次优化都砍一半时间,但数据量涨 6 倍,时间还是不够。数据库的极限在哪?”
大师:“我们来看看四种写入方式的极限吞吐——”
# 方式 1:ORM add 分批 flush → 约 500-1000 行/秒
# 方式 2:bulk_insert_mappings → 约 5000-8000 行/秒
# 方式 3:session.execute(insert().values([...])) → 约 20000-50000 行/秒
# 方式 4:COPY (PostgreSQL 原生) → 约 100000-500000 行/秒
小胖:“差了 1000 倍!方式 4 的 COPY 为什么这么快?”
大师:“COPY 是 PostgreSQL 的原生批量导入命令——它绕过了 SQL 解析、计划优化、结果集处理的全套流程,直接以二进制流写入。这就像快递站直接从车屁股卸货——不需要每个包裹都走一遍安检和分拣。”
小胖:“技术映射:ORM add = 逐个排队寄信;Core INSERT = 批量装麻袋寄;COPY = 整车厢式卡车专线直达。”
小白:“那读写分离呢?我们这个导入作业在写入时,报表查询卡得不行——能不能把查询分离出去?”
大师:"读写分离有两种实现方式:
- 应用层路由:
RoutingSession根据操作类型(读/写)选择不同的 Engine。 - 数据库层复制:主 PostgreSQL (writes) → 流复制 → 只读副本 (reads)。"
# 应用层路由
class RoutingSession(Session):
def get_bind(self, mapper=None, clause=None, **kw):
if self._is_reading(clause):
return read_engine # 连接只读副本
return write_engine # 连接主库
小白:“技术映射:读写分离 = 高速公路双向车道——写入走主路(主库),查询走辅路(只读副本),互不干扰。”
大师:“但要小心’写后读’的一致性延迟——主库写入后副本有复制延迟(可能 50-200ms)。如果你的业务是’写入后立即查询’——必须从主库读。”
小胖:“那 OLTP 和 OLAP 的 Session 需要分开吗?”
大师:“强烈建议分开——至少两个独立的 SessionFactory。OLTP Session 关闭 expire_on_commit、关闭 autoflush、使用 READ COMMITTED 隔离级别。OLAP(报表)Session 使用专用连接池(pool_size 更小但连接更稳定),使用 REPEATABLE READ 以保证报表查询的一致性。”
三、项目实战
实战目标
对比四种写入方式的吞吐极限,实现应用层读写分离路由,设计 OLTP/OLAP Session 隔离方案,给出容量规划建议。
步骤一:四种写入方式性能基准
"""ch35_extreme_performance.py —— 极端性能实战"""
from sqlalchemy import (
create_engine, String, Integer, Numeric, DateTime,
ForeignKey, text, func, select, insert, event,
)
from sqlalchemy.orm import (
DeclarativeBase, Mapped, mapped_column, relationship,
Session, sessionmaker,
)
from datetime import datetime
import time, random, io, csv
# 写库(主库)
write_engine = create_engine(
"postgresql+psycopg://nebula:nebula_dev@localhost:5432/order_center",
echo=False,
pool_size=5,
max_overflow=10,
)
# 读库(只读副本——本示例中使用同一库,实际部署中分离)
read_engine = create_engine(
"postgresql+psycopg://nebula:nebula_dev@localhost:5432/order_center",
echo=False,
pool_size=3,
max_overflow=5,
)
class Base(DeclarativeBase):
pass
class OrderItem(Base):
__tablename__ = "perf_import_items"
id: Mapped[int] = mapped_column(primary_key=True)
order_id: Mapped[int] = mapped_column(Integer, index=True)
product_name: Mapped[str] = mapped_column(String(200))
unit_price: Mapped[float] = mapped_column(Numeric(12, 2))
quantity: Mapped[int] = mapped_column(Integer)
imported_at: Mapped[datetime] = mapped_column(DateTime, default=func.now())
Base.metadata.drop_all(write_engine)
Base.metadata.create_all(write_engine)
# =============================================
# 基准测试:四种写入方式
# =============================================
NUM_ROWS = 50_000 # 5 万行作为基准
def generate_data(n: int):
return [
{
"order_id": random.randint(1, 10000),
"product_name": f"商品-{random.randint(1, 5000)}",
"unit_price": round(random.uniform(1, 1000), 2),
"quantity": random.randint(1, 100),
}
for _ in range(n)
]
def benchmark_import(label, fn):
"""执行并计时"""
# 清空表
with write_engine.begin() as conn:
conn.execute(text("DELETE FROM perf_import_items"))
conn.commit()
start = time.perf_counter()
rows = fn()
elapsed = time.perf_counter() - start
rate = rows / elapsed if elapsed > 0 else 0
print(f" [{label}] 耗时: {elapsed:.1f}s | 速率: {rate:,.0f} 行/秒")
return elapsed, rate
print("=== 四种写入方式性能基准 ===")
# 方式 1:ORM add_all + 分批 flush
def orm_add_batch():
data = generate_data(NUM_ROWS)
Factory = sessionmaker(bind=write_engine)
with Factory() as s:
batch_size = 1000
for i in range(0, NUM_ROWS, batch_size):
batch = [OrderItem(**d) for d in data[i:i+batch_size]]
s.add_all(batch)
s.flush()
s.commit()
return NUM_ROWS
benchmark_import("ORM add_all 分批1000", orm_add_batch)
# 方式 2:bulk_insert_mappings
def orm_bulk():
data = generate_data(NUM_ROWS)
Factory = sessionmaker(bind=write_engine)
with Factory() as s:
s.bulk_insert_mappings(OrderItem, data)
s.commit()
return NUM_ROWS
benchmark_import("bulk_insert_mappings", orm_bulk)
# 方式 3:Core insert (多值)
def core_insert_multi():
data = generate_data(NUM_ROWS)
with write_engine.begin() as conn:
conn.execute(insert(OrderItem.__table__), data)
return NUM_ROWS
benchmark_import("Core INSERT 多值", core_insert_multi)
# 方式 4:COPY (PostgreSQL 原生)
def copy_from_csv():
data = generate_data(NUM_ROWS)
# 将数据写入内存 CSV
output = io.StringIO()
writer = csv.writer(output)
for row in data:
writer.writerow([
row["order_id"], row["product_name"],
row["unit_price"], row["quantity"],
datetime.now().isoformat(),
])
output.seek(0)
raw_conn = write_engine.raw_connection()
try:
cursor = raw_conn.cursor()
cursor.copy_from(
output,
"perf_import_items",
sep=",",
columns=("order_id", "product_name", "unit_price", "quantity", "imported_at"),
)
raw_conn.commit()
finally:
raw_conn.close()
return NUM_ROWS
benchmark_import("COPY (PostgreSQL)", copy_from_csv)
步骤二:读写分离路由
# =============================================
# 读写分离路由
# =============================================
from sqlalchemy.orm import Session as _Session
from sqlalchemy.sql import Select
class RoutingSession(_Session):
"""根据 SQL 操作类型路由到不同的数据库连接"""
_is_write = False # 追踪当前操作是否为写操作
def get_bind(self, mapper=None, clause=None, **kw):
"""决定使用哪个数据库引擎"""
# INSERT/UPDATE/DELETE → 主库(write_engine)
if self._is_write or (clause is not None and not isinstance(clause, Select)):
return write_engine
# SELECT → 只读副本(read_engine)
return read_engine
def execute(self, statement, *args, **kwargs):
# 判断是否为写操作
from sqlalchemy.sql import Update, Delete, Insert
self._is_write = isinstance(statement, (Update, Delete, Insert))
return super().execute(statement, *args, **kwargs)
RoutingFactory = sessionmaker(class_=RoutingSession)
print("\n=== 读写分离测试 ===")
with RoutingFactory() as s:
# SELECT 会路由到 read_engine
item = s.execute(
select(OrderItem).where(OrderItem.id == 1)
).scalars().first()
print(f" [READ] → read_engine: {item.product_name if item else 'N/A'}")
with RoutingFactory() as s:
# INSERT 会路由到 write_engine
new_item = OrderItem(
order_id=9999, product_name="测试商品",
unit_price=10, quantity=1,
)
s.add(new_item)
s.commit()
print(" [WRITE] → write_engine: INSERT 完成")
print("\n 注意事项:")
print(" 1. 写后读可能遇到复制延迟(read_engine 未同步)")
print(" 2. 对于需要'写入后立即读取'的操作,使用显式的 write_engine")
print(" 3. 事务中的混合读写在同一个 bind 上执行")
步骤三:OLTP / OLAP Session 隔离
# =============================================
# OLTP / OLAP Session 隔离
# =============================================
# OLTP Session Factory——面向在线事务
OLTPFactory = sessionmaker(
bind=write_engine,
expire_on_commit=False, # 避免 commit 后属性访问触发 SELECT
autoflush=True, # 查询前自动 flush pending 变更
)
# OLAP Session Factory——面向分析查询
OLAPFactory = sessionmaker(
bind=read_engine, # 使用只读副本
expire_on_commit=False,
autoflush=False, # 分析查询不需要 autoflush
# 使用 REPEATABLE READ 保证同一事务内数据一致性
)
print("\n=== OLTP / OLAP Session 隔离 ===")
# OLTP 操作:快速写入
with OLTPFactory() as s:
item = OrderItem(
order_id=9998, product_name="OLTP测试",
unit_price=50, quantity=2,
)
s.add(item)
s.commit()
print(f" [OLTP] 写入完成: id={item.id}")
# OLAP 操作:只读分析(使用独立连接池)
with OLAPFactory() as s:
result = s.execute(
select(
func.count(OrderItem.id).label("total"),
func.avg(OrderItem.unit_price).label("avg_price"),
func.sum(OrderItem.quantity).label("total_qty"),
)
).first()
print(f" [OLAP] 统计: 总数={result.total}, 均价={result.avg_price:.2f}, 总数量={result.total_qty}")
# 总结
print("\n OLTP Session 特点:")
print(" - 连接主库 (write_engine)")
print(" - expire_on_commit=False(避免隐式 SELECT)")
print(" - autoflush=True(确保写入顺序正确)")
print(" - 事务隔离级别: READ COMMITTED")
print(" OLAP Session 特点:")
print(" - 连接只读副本 (read_engine)")
print(" - expire_on_commit=False")
print(" - autoflush=False(禁止写入)")
print(" - 事务隔离级别: REPEATABLE READ")
print(" - 独立连接池(不与 OLTP 争抢连接)")
步骤四:批量导入的生产优化
# =============================================
# 生产优化技巧汇总
# =============================================
PROD_OPTIMIZATION = """
生产级批量导入优化清单:
1. 关闭 autovacuum(导入前):
ALTER TABLE perf_import_items SET (autovacuum_enabled = false);
-- 导入完成后重新打开并手动 VACUUM ANALYZE
2. 删除索引(导入前):
DROP INDEX IF EXISTS idx_perf_import_items_order_id;
-- 导入后重建索引:CREATE INDEX CONCURRENTLY ...
3. 增大 work_mem 和 maintenance_work_mem(会话级):
SET work_mem = '256MB';
SET maintenance_work_mem = '512MB';
4. 使用 UNLOGGED 表(可丢失数据但极快):
CREATE UNLOGGED TABLE import_staging (...) ;
5. 多 Worker 并行写入——每个 Worker 使用独立连接:
- 每个 Worker 从独立的数据分片读取
- 使用独立的 Engine(避免连接池争抢)
- 最后用 INSERT INTO ... SELECT 合并到主表
6. 关闭 WAL(wal_level=minimal + 批量导入):
-- 危险操作——只能在初始导入时使用
-- 正常生产不能关 WAL
7. 使用 partition(表分区):
- 按日期分区 ORDER_ITEMS_202401, ORDER_ITEMS_202402...
- 导入作业只写当天分区(减小表大小、加速 VACUUM)
"""
print("\n=== 生产优化技巧 ===")
print(PROD_OPTIMIZATION)
# =============================================
# 容量规划公式
# =============================================
CAPACITY_PLANNING = """
容量规划建议:
连接数需求 = (API Worker 数 × pool_size) + (定时任务 × pool_size) + 预留(20%)
= (30 × 5) + (5 × 3) + 20% ≈ 200 连接
数据库 max_connections = 连接数需求 × 1.5 ≈ 300
磁盘 I/O 估算:
每秒写入 MB = 写入速率(行/秒) × 行大小(MB)
300 万行 × 200 bytes/行 = 600MB
目标 5 分钟 → 120MB/min → 2MB/s 持续写入
峰值 I/O = 写入 + WAL + 索引更新 ≈ 2MB/s × 4 = 8MB/s
内存估算:
连接数 × work_mem (4MB) + shared_buffers (建议 RAM 的 25%)
200 × 4MB + 4GB = 4.8GB
"""
print("\n=== 容量规划 ===")
print(CAPACITY_PLANNING)
可能遇到的坑及解决方法
- COPY 命令在事务中可能无法与其他操作共享连接
- 现象:
psycopg.errors.ActiveSqlTransaction: COPY ... cannot run inside a transaction block。 - 根因:部分 PostgreSQL 配置不允许在事务块中执行 COPY。
- 解决:使用
raw_connection()获取底层 DBAPI 连接,在 autocommit 模式下执行 COPY。
- 读写分离中"写后读"看到旧数据
- 现象:写入订单后立即查询——查询结果中没有新订单。
- 根因:读请求被路由到了只读副本,而副本有复制延迟(50ms-几秒)。
- 解决:对"写入后立即读取"的操作,临时将查询也路由到主库(
using(write_engine))。
- 多 Worker 并行写入时连接池耗尽
- 现象:10 个 Worker 共享 pool_size=5 的引擎——5 个等待超时。
- 解决:增大 pool_size 或使用 NullPool(每个 Worker 独立连接)。或使用 PgBouncer 做连接池前置。
session.execute(insert().values([...]))对 Python 内存压力大
- 现象:300 万条数据
values([...])→ 列表里有 300 万个字典 → Python 进程 OOM。 - 解决:分批 values(每批 5000-10000 行),或使用
add_all()+flush(),或直接用 COPY。
测试验证
# tests/test_ch35_performance.py
import pytest
from sqlalchemy import create_engine, text, insert, String, Integer, Numeric
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, sessionmaker, Session
import time
def test_core_insert_faster_than_orm_bulk():
"""验证 Core INSERT 比 ORM bulk_insert 更快"""
engine = create_engine("sqlite:///:memory:", echo=False)
class Base(DeclarativeBase):
pass
class T(Base):
__tablename__ = "t"
id: Mapped[int] = mapped_column(primary_key=True)
name: Mapped[str] = mapped_column(String(50))
Base.metadata.create_all(engine)
data = [{"name": f"item-{i}"} for i in range(5000)]
# ORM bulk
F = sessionmaker(bind=engine)
t0 = time.perf_counter()
with F() as s:
s.bulk_insert_mappings(T, data)
s.commit()
bulk_time = time.perf_counter() - t0
# 清表
with engine.begin() as conn:
conn.execute(text("DELETE FROM t"))
# Core INSERT
t1 = time.perf_counter()
with engine.begin() as conn:
conn.execute(insert(T.__table__), data)
core_time = time.perf_counter() - t1
assert core_time <= bulk_time * 1.2, f"Core INSERT 应与 bulk 相当或更快"
def test_routing_session_routes_read_and_write():
"""验证 RoutingSession 正确路由读/写"""
engine_r = create_engine("sqlite:///:memory:", echo=False)
engine_w = create_engine("sqlite:///:memory:", echo=False)
class Base(DeclarativeBase):
pass
class T(Base):
__tablename__ = "rt"
id: Mapped[int] = mapped_column(primary_key=True)
name: Mapped[str] = mapped_column(String(50))
Base.metadata.create_all(engine_w)
Base.metadata.create_all(engine_r)
from sqlalchemy import select, insert
from sqlalchemy.sql import Select
class Router(Session):
def get_bind(self, mapper=None, clause=None, **kw):
if clause is not None and not isinstance(clause, Select):
return engine_w
return engine_r
RF = sessionmaker(class_=Router)
with RF() as s:
# 写入
s.execute(insert(T.__table__).values(name="test"))
assert s.get_bind().url.database == engine_w.url.database
def test_oltp_session_no_expire_on_commit():
"""验证 OLTP Session 不 expire 已提交的属性"""
engine = create_engine("sqlite:///:memory:", echo=False)
class Base(DeclarativeBase):
pass
class T(Base):
__tablename__ = "oltp"
id: Mapped[int] = mapped_column(primary_key=True)
name: Mapped[str] = mapped_column(String(50))
Base.metadata.create_all(engine)
F = sessionmaker(bind=engine, expire_on_commit=False)
with F() as s:
t = T(name="oltp_test")
s.add(t)
s.commit()
# expire_on_commit=False → commit 后可以直接读取属性
assert t.name == "oltp_test", "expire_on_commit=False 时属性应保持可用"
四、项目总结
四种写入方式对比
| 方式 | 吞吐 (行/秒) | Python 开销 | 事务安全 | 适用数据量 |
|---|---|---|---|---|
| ORM add_all + flush | 500-2000 | 极高 | 是 | < 1万 |
| bulk_insert_mappings | 5000-10000 | 中 | 是 | 1万-5万 |
| Core INSERT 多值 | 20000-80000 | 低 | 是 | 5万-100万 |
| COPY (PostgreSQL) | 100000-500000 | 极低 | 是 | > 100万 |
适用场景
- 业务事务(ORM)——下单、退款——数据量小,需要对象导航。
- 定时批量导入(Core INSERT / COPY)——商品价格同步、物流数据导入。
- 报表统计(只读 Core / OLAP Session)——聚合查询走只读副本。
- ETL 数据清洗(COPY + 数据库函数)——利用 PostgreSQL 的 COPY 避免 Python 参与。
不适用场景:
- 在线服务中直接使用 COPY——COPY 需要独占总线,可能阻塞正常 API 查询。
- 数据量 < 10000 行——使用 ORM add 的便利性比 COPY 的性能提升更有价值。
注意事项
COPY绕过 SQL 解析——触发器不会触发(如果需要触发器,使用INSERT替代)。- 读写分离的复制延迟——写入后立即读取需要路由到主库。
- 批量导入时监控连接池——多 Worker + COPY 可能耗尽连接。
- OLAP 查询不要与 OLTP 共用连接池——长查询占用连接导致 OLTP 排队超时。
常见踩坑经验
案例 1:COPY 后主键序列(Sequence)未更新
- 现象:COPY 1 万条数据后,用 ORM 插入新数据报 duplicate key violation。
- 根因:COPY 不更新 PostgreSQL 的
SEQUENCE——自增序列值仍停在 COPY 之前的数字。 - 修复:
SELECT setval('table_id_seq', (SELECT MAX(id) FROM table));
案例 2:读写分离导致事务中的一致性视图错乱
- 现象:事务内先 UPDATE 然后 SELECT——UPDATE 在主库,SELECT 在副本——看不到刚更新的数据。
- 修复:事务内所有操作应在同一个 bind 上执行——修改
RoutingSession在事务内固定 bind。
案例 3:COPY 导入中的 CSV 格式错误难以定位
- 现象:第 45,678 行 CSV 数据格式错误导致 COPY 失败——但错误信息只显示"COPY error",不显示行号。
- 解决:导入前用
csv模块验证格式。将数据拆成小份(每份 1 万行)分别 COPY,缩小出错范围。
思考题
-
PostgreSQL 的
COPY命令可以读取本地文件:COPY table FROM '/path/to/data.csv'。但云数据库(如 RDS)通常限制对服务器文件系统的访问。如何在不使用服务器文件系统的条件下,达到 COPY 级别的导入性能?cursor.copy_expert()的STDIN模式能否解决这个问题? -
如果需要对导入的百万行数据进行逐行处理(如根据现有数据决定 INSERT 或 UPDATE),纯 COPY 无法满足(COPY 只在表级操作)。你怎么在"批量处理"和"COPY 性能"之间找到平衡?INSERT … ON CONFLICT DO UPDATE 和 COPY + MERGE 各有什么优劣?
延伸阅读与资源
NumPy 从入门到生产落地:全链路实战指南(科学计算/向量化)
Redis 8 实战精讲:从 CRUD 到源码,构建高可用缓存系统
Redis 实战修炼与原理进阶
Python 3实战精进:从脚本到高并发订单引擎
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
MongoDB 实战进阶与内核修炼
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析
更多推荐




所有评论(0)