一、项目背景

"库存变成了负数!"2024 年双十一凌晨 0 点 5 分,星云电商监控面板爆红。

那场事故的起因看似简单:一款热门耳机库存只剩 2 件,但同一秒内涌入了 5 个购买请求。由于下单接口的事务隔离级别配置存在问题,5 个请求几乎同时读到库存为 2,各自判断"库存足够",然后分别执行了 UPDATE products SET inventory = inventory - 1——结果库存从 2 径直掉到了 -3。

更令人头皮发麻的是后续的连锁反应:超卖触发了库存为负的检查约束(CHECK inventory >= 0),数据库直接报错。但下单服务已经创建了订单记录(在同一个事务内),有三个用户的订单状态变成了"已支付但库存不存在"的幽灵订单。这直接导致了 3 笔退款、2 个投诉和 1 个客诉升级。

事后复盘,事故的三个根因是:

  1. 事务隔离级别使用了默认的 READ COMMITTED——在"读-判断-写"场景下,多个并发事务可能读到相同的"旧值",导致 lost update。
  2. 没有使用乐观并发控制(版本号)——当两个事务同时修改同一行时,没有机制让第二个事务感知到"数据已被修改"。
  3. 事务粒度过大——下单接口将"查库存→创建订单→扣减库存→发送通知"全部放在一个事务里,持锁时间过长。

本章将围绕事务的四个核心概念展开:事务边界(Connection 级 vs Session 级)、隔离级别(READ COMMITTED / REPEATABLE READ / SERIALIZABLE)、乐观并发控制(version_id_col)、以及失败重试与幂等设计。

二、项目设计

场景:周三下午,事故复盘会。小胖一脸沉重,大师带着白板走进会议室。

大师:“小胖,你来说说,下单接口当时发生了什么?”

小胖:“5 个请求同时进来,每个请求做了三件事:第一步查库存(SELECT inventory FROM products WHERE id = 1),第二步判断库存是否 > 0,第三步 UPDATE 扣减库存。问题是——5 个请求都读到 inventory = 2,然后都认为可以下单……”

大师:“这叫 Lost Update(丢失更新)——在 READ COMMITTED 隔离级别下,一个事务的 SELECT 看不到其他事务未提交的变更。但两个并发事务读到同样的初始值,各自基于这个值修改,最终只有一个事务的修改生效,另一个被覆盖。”

小白:“那隔离级别是什么?能不能解决这个问题?”

大师:“隔离级别决定了一个事务在多大程度上’看到’其他并发事务的变更。从上到下四个级别:”

隔离级别 脏读 不可重复读 幻读 丢失更新
READ UNCOMMITTED 可能 可能 可能 可能
READ COMMITTED (PG默认) 可能 可能 可能
REPEATABLE READ 可能(PG不) 可能
SERIALIZABLE

小胖:“技术映射:隔离级别 = 房间的透明度和隔音程度。READ UNCOMMITTED 是全透明的玻璃房,SERIALIZABLE 是不透明隔音房。但 SERIALIZABLE 是不是很慢?”

大师:“性能确实有影响,但 PostgreSQL 的 SERIALIZABLE 通过 SSI(可序列化快照隔离)做了优化,不像传统方案全程加锁。对于扣库存这种强一致性需求,有两种路径:一是提升隔离级别到 REPEATABLE READ + SELECT FOR UPDATE;二是使用乐观并发控制——版本号。”

小白:“技术映射:乐观并发 = 先做再说,提交时检查有没有人抢先改了。版本号怎么用?”

大师

class Product(Base):
    __tablename__ = "products"
    id: Mapped[int] = mapped_column(primary_key=True)
    inventory: Mapped[int] = mapped_column(Integer)
    version_id: Mapped[int] = mapped_column(
        Integer, nullable=False, server_default=text("1")
    )

    __mapper_args__ = {"version_id_col": version_id}

大师:“当你配置了 version_id_col,每次 UPDATE 时,SQLAlchemy 会自动在 WHERE 子句中加上 AND version_id = :当前版本号,同时把版本号 +1。如果另一个事务已经先改了同一行(版本号已经不匹配),UPDATE 会返回 rowcount = 0,SQLAlchemy 抛出 StaleDataError。”

-- SQLAlchemy 自动生成的 SQL
UPDATE products SET inventory = 1, version_id = 3
WHERE id = 1 AND version_id = 2  -- 版本号不匹配时 rowcount=0

小胖:“技术映射:version_id = 并发修改的检测哨兵。这就相当于去银行办业务,柜员会先看你的排队号码,如果号码已经过号了,就拒绝服务。”

小白:“那事务边界呢?Connection 级和 Session 级的事务有什么区别?”

大师:“Connection 级事务是数据库原生事务——conn.begin() 相当于 BEGIN SQL 命令。Session 级事务是 ORM 的封装——但你 session.rollback() 时,它底层调用的是 conn.rollback()。所以 Session 的事务本质就是 Connection 事务的代理。”

小胖:“那我什么时候应该用 ORM Session,什么时候应该在 Connection 上直接管理事务?”

大师:“如果你只是 Core 操作(批量写、纯 SQL),用 Connection 级事务更轻量。如果你涉及 ORM 对象的增删改(Session 的 flush/commit),用 Session 级事务。但核心原则是——事务粒度越小越好,持锁时间越短越好。”

小白:“技术映射:事务粒度 = 结账的批量大小——太大排队长,太小柜台效率低。”

三、项目实战

实战目标

模拟"扣库存下单"的完整事务流程,涵盖:事务边界管理、乐观并发控制(版本号)、失败重试与幂等机制。对比无版本号 vs 有版本号在高并发下的差异。

步骤一:模型声明(带版本号的 Product)

"""ch13_transaction_isolation.py —— 事务、隔离级别、乐观并发"""

from sqlalchemy import (
    create_engine, String, Integer, Numeric, DateTime,
    ForeignKey, text, func, select, update, and_, or_, event,
)
from sqlalchemy.orm import (
    DeclarativeBase, Mapped, mapped_column, relationship,
    Session, sessionmaker,
)
from sqlalchemy.exc import StaleDataError
from datetime import datetime
from typing import Optional
import time, threading

engine = create_engine(
    "postgresql+psycopg://nebula:nebula_dev@localhost:5432/order_center",
    echo=False,
    # 默认隔离级别 READ COMMITTED
    # 可在 create_engine 时指定:
    # isolation_level="REPEATABLE READ" 或 "SERIALIZABLE"
)

class Base(DeclarativeBase):
    pass

class Product(Base):
    __tablename__ = "txn_products"
    id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
    sku: Mapped[str] = mapped_column(String(30), unique=True, nullable=False)
    title: Mapped[str] = mapped_column(String(200), nullable=False)
    unit_price: Mapped[float] = mapped_column(Numeric(12, 2), nullable=False)
    inventory: Mapped[int] = mapped_column(Integer, nullable=False, server_default=text("0"))
    # 版本号——乐观并发控制
    version_id: Mapped[int] = mapped_column(Integer, nullable=False, server_default=text("1"))

    __mapper_args__ = {"version_id_col": version_id}
    __table_args__ = (CheckConstraint("inventory >= 0", name="ck_inventory"),)

class Order(Base):
    __tablename__ = "txn_orders"
    id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
    order_no: Mapped[str] = mapped_column(String(32), unique=True, nullable=False)
    user_name: Mapped[str] = mapped_column(String(50), nullable=False)
    product_id: Mapped[int] = mapped_column(ForeignKey("txn_products.id"), nullable=False)
    quantity: Mapped[int] = mapped_column(Integer, nullable=False)
    total_amount: Mapped[float] = mapped_column(Numeric(12, 2), nullable=False)
    status: Mapped[str] = mapped_column(String(20), server_default=text("'pending'"))
    created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now())

    __table_args__ = (
        UniqueConstraint("order_no", name="uq_order_no"),
        Index("ix_txn_orders_user", "user_name"),
    )

Base.metadata.create_all(engine)
SessionFactory = sessionmaker(bind=engine, autocommit=False, autoflush=False)

步骤二:基本事务——扣库存+下单(单线程正确版)

# =============================================
# 基本事务:扣库存 + 创建订单
# =============================================

def place_order_optimistic(
    user_name: str, product_id: int, quantity: int, max_retries: int = 3
) -> Optional[Order]:
    """
    乐观并发控制的扣库存下单
    如果版本号冲突(并发修改),自动重试
    """
    for attempt in range(1, max_retries + 1):
        session = SessionFactory()
        try:
            # 1. 查询产品(带版本号)
            product = session.get(Product, product_id)
            if not product:
                raise ValueError(f"产品 {product_id} 不存在")

            # 2. 业务校验
            if product.inventory < quantity:
                raise ValueError(f"库存不足: 需要 {quantity}, 剩余 {product.inventory}")

            # 3. 扣减库存(直接修改 persistent 对象属性)
            product.inventory -= quantity
            # SQLAlchemy 自动将 version_id 加 1,flush 时生成:
            # UPDATE ... SET inventory = ..., version_id = version_id + 1
            # WHERE id = ? AND version_id = ?

            # 4. 创建订单
            import random
            order = Order(
                order_no=f"ORD-{datetime.now().strftime('%Y%m%d%H%M%S')}-{random.randint(1000, 9999)}",
                user_name=user_name,
                product_id=product_id,
                quantity=quantity,
                total_amount=product.unit_price * quantity,
                status="paid",
            )
            session.add(order)
            session.commit()

            print(f"  [attempt {attempt}] 下单成功: {order.order_no}, "
                  f"库存剩余: {product.inventory}, 版本: {product.version_id}")
            return order

        except StaleDataError:
            # 版本号冲突——数据已被其他事务修改,当前事务回滚并重试
            session.rollback()
            print(f"  [attempt {attempt}] 版本冲突(并发修改),即将重试...")
            time.sleep(0.01 * attempt)  # 退避策略
            continue
        except Exception as e:
            session.rollback()
            print(f"  [attempt {attempt}] 错误: {e}")
            return None
        finally:
            session.close()

    print(f"  重试 {max_retries} 次后仍失败")
    return None

# 测试:插入基础数据
print("=== 基本事务:扣库存 + 下单 ===")
with SessionFactory() as session:
    # 确保数据存在
    existing = session.execute(select(Product).where(Product.sku == "SKU-HOT-001")).scalars().first()
    if not existing:
        product = Product(
            sku="SKU-HOT-001", title="热门耳机 Pro", unit_price=299.00, inventory=10
        )
        session.add(product)
        session.commit()
        product_id = product.id
        print(f"初始数据: {product.title}, 库存={product.inventory}, 版本={product.version_id}")
    else:
        product_id = existing.id
        print(f"已有数据: id={product_id}, 库存={existing.inventory}")

# 单线程测试
order = place_order_optimistic("小胖", product_id=product_id, quantity=2)
print()

步骤三:高并发模拟——对比无版本号 vs 有版本号

# =============================================
# 并发模拟:多线程同时下单
# =============================================

def place_order_without_version(user_name: str, product_id: int, quantity: int):
    """无版本号的下单——存在 lost update 风险"""
    session = SessionFactory()
    try:
        product = session.get(Product, product_id)
        if product.inventory < quantity:
            session.rollback()
            return None
        product.inventory -= quantity
        order = Order(
            order_no=f"NV-{threading.get_ident()}-{int(time.time()*1000)}",
            user_name=user_name, product_id=product_id,
            quantity=quantity, total_amount=product.unit_price * quantity,
            status="paid",
        )
        session.add(order)
        session.commit()
        return order
    except Exception:
        session.rollback()
        return None
    finally:
        session.close()

def concurrent_test():
    """多线程并发下单测试"""
    with SessionFactory() as session:
        # 重置库存
        session.execute(
            update(Product).where(Product.sku == "SKU-HOT-002").values(inventory=5, version_id=1)
        )
        product = Product(sku="SKU-HOT-002", title="测试商品", unit_price=50.00, inventory=5) \
            if not session.execute(select(Product).where(Product.sku == "SKU-HOT-002")).scalars().first() \
            else None
        if product:
            session.add(product)
        session.commit()
        product_id = session.execute(select(Product.id).where(Product.sku == "SKU-HOT-002")).scalar()

    print(f"并发测试: 库存=5, 10 个线程各抢 1 件")

    results = []
    errors = []

    def thread_task(user_idx, use_version=True):
        try:
            if use_version:
                order = place_order_optimistic(f"user_{user_idx}", product_id, 1, max_retries=5)
            else:
                order = place_order_without_version(f"user_{user_idx}", product_id, 1)
            results.append((user_idx, order is not None))
        except Exception as e:
            errors.append((user_idx, str(e)))

    threads = []
    for i in range(10):
        t = threading.Thread(target=thread_task, args=(i, True))
        threads.append(t)
        t.start()

    for t in threads:
        t.join()

    success = sum(1 for _, ok in results if ok)
    fail = sum(1 for _, ok in results if not ok)
    print(f"  成功: {success}, 失败: {fail}")
    if errors:
        print(f"  错误: {errors[:3]}")

    with SessionFactory() as session:
        product = session.get(Product, product_id)
        print(f"  最终库存: {product.inventory} (预期 0), 版本: {product.version_id}")
        if product.inventory < 0:
            print(f"  [严重] 库存为负——发生了超卖!")
    print()

print("=== 高并发下单测试 ===")
concurrent_test()

步骤四:事务隔离级别对比实验

# =============================================
# 隔离级别对比实验
# =============================================

def test_isolation_levels():
    """演示不同隔离级别下脏读/不可重复读的行为"""
    from sqlalchemy import create_engine as ce

    print("=== 隔离级别对比实验 ===")

    # 创建两个引擎:一个默认 READ COMMITTED,一个 REPEATABLE READ
    eng_rc = ce("postgresql+psycopg://nebula:nebula_dev@localhost:5432/order_center", echo=False)
    eng_rr = ce("postgresql+psycopg://nebula:nebula_dev@localhost:5432/order_center",
                echo=False, isolation_level="REPEATABLE READ")

    # 准备数据
    with Session(eng_rc) as s:
        s.execute(text("CREATE TEMP TABLE IF NOT EXISTS _iso_test (id INT, val INT)"))
        s.execute(text("DELETE FROM _iso_test"))
        s.execute(text("INSERT INTO _iso_test VALUES (1, 100)"))
        s.commit()

    # 实验 1:READ COMMITTED 下的不可重复读
    print("--- READ COMMITTED ---")
    s1 = Session(eng_rc)
    s2 = Session(eng_rc)

    row1 = s1.execute(text("SELECT val FROM _iso_test WHERE id = 1")).scalar()
    print(f"  事务A 第一次读: val = {row1}")

    s2.execute(text("UPDATE _iso_test SET val = 200 WHERE id = 1"))
    s2.commit()
    print(f"  事务B 修改 val 为 200 并提交")

    row2 = s1.execute(text("SELECT val FROM _iso_test WHERE id = 1")).scalar()
    print(f"  事务A 第二次读: val = {row2} (已变成 200——这就是不可重复读)")
    s1.close()
    s2.close()

    # 实验 2:REPEATABLE READ 防止不可重复读
    print("\n--- REPEATABLE READ ---")
    with Session(eng_rc) as s:
        s.execute(text("UPDATE _iso_test SET val = 100 WHERE id = 1"))
        s.commit()

    s1 = Session(eng_rr)
    s2 = Session(eng_rc)

    row1 = s1.execute(text("SELECT val FROM _iso_test WHERE id = 1")).scalar()
    print(f"  事务A 第一次读: val = {row1}")

    s2.execute(text("UPDATE _iso_test SET val = 200 WHERE id = 1"))
    s2.commit()
    print(f"  事务B 修改 val 为 200 并提交")

    row2 = s1.execute(text("SELECT val FROM _iso_test WHERE id = 1")).scalar()
    print(f"  事务A 第二次读: val = {row2} (仍然是 100——REPEATABLE READ 快照隔离)")
    s1.close()
    s2.close()

test_isolation_levels()

步骤五:幂等设计——防止重复下单

# =============================================
# 幂等设计:重复提交保护
# =============================================

def idempotent_place_order(
    user_name: str, product_id: int, quantity: int,
    idempotency_key: str,  # 幂等键(前端生成,如 UUID)
) -> Optional[Order]:
    """
    幂等下单:同一 idempotency_key 重复调用只创建一次订单
    """
    session = SessionFactory()
    try:
        # 1. 幂等检查:是否已有相同幂等键的订单
        existing = session.execute(
            select(Order).where(Order.order_no == idempotency_key)
        ).scalars().first()

        if existing:
            print(f"  检测到重复请求(幂等键: {idempotency_key}),返回已有订单")
            return existing

        # 2. 正常下单流程(带版本号乐观锁)
        product = session.get(Product, product_id)
        if not product or product.inventory < quantity:
            raise ValueError("库存不足")

        product.inventory -= quantity

        order = Order(
            order_no=idempotency_key,  # 把幂等键作为订单号
            user_name=user_name,
            product_id=product_id,
            quantity=quantity,
            total_amount=product.unit_price * quantity,
            status="paid",
        )
        session.add(order)
        session.commit()
        return order

    except StaleDataError:
        session.rollback()
        # 版本冲突时的处理:让调用方重试
        return idempotent_place_order(user_name, product_id, quantity, idempotency_key)

    except Exception as e:
        session.rollback()
        raise
    finally:
        session.close()

print("=== 幂等下单测试 ===")
import uuid
idem_key = f"IDEM-{uuid.uuid4()}"
o1 = idempotent_place_order("小白", product_id=1, quantity=1, idempotency_key=idem_key)
o2 = idempotent_place_order("小白", product_id=1, quantity=1, idempotency_key=idem_key)
print(f"  两次调用返回同一订单: {o1 is o2 or (o1 and o2 and o1.id == o2.id)}")

完整代码清单

"""ch13_transaction_complete.py —— 事务与一致性完整示例"""

from sqlalchemy import create_engine, String, Integer, Numeric, DateTime, ForeignKey, text, func, select, update
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, sessionmaker, Session
from sqlalchemy.exc import StaleDataError
from datetime import datetime
from order_center.config import DATABASE_URL
import time

engine = create_engine(DATABASE_URL, echo=False)

class Base(DeclarativeBase):
    pass

class Product(Base):
    __tablename__ = "tx_products"
    id: Mapped[int] = mapped_column(primary_key=True)
    sku: Mapped[str] = mapped_column(String(30), unique=True, nullable=False)
    title: Mapped[str] = mapped_column(String(200), nullable=False)
    unit_price: Mapped[float] = mapped_column(Numeric(12, 2), nullable=False)
    inventory: Mapped[int] = mapped_column(Integer, nullable=False)
    version_id: Mapped[int] = mapped_column(Integer, nullable=False, server_default=text("1"))
    __mapper_args__ = {"version_id_col": version_id}

Base.metadata.create_all(engine)
Factory = sessionmaker(bind=engine)

def deduct_inventory_optimistic(session: Session, product_id: int, quantity: int) -> Product:
    """乐观锁扣减库存——版本号冲突时让调用方重试"""
    product = session.get(Product, product_id)
    if not product or product.inventory < quantity:
        raise ValueError("库存不足")
    product.inventory -= quantity
    return product

def retry_on_stale(fn, max_retries=3, *args, **kwargs):
    """自动重试 StaleDataError"""
    for i in range(max_retries):
        session = Factory()
        try:
            result = fn(session, *args, **kwargs)
            session.commit()
            return result
        except StaleDataError:
            session.rollback()
            if i == max_retries - 1:
                raise
            time.sleep(0.01 * (i + 1))
        except Exception:
            session.rollback()
            raise
        finally:
            session.close()

if __name__ == "__main__":
    print("事务模块就绪")

可能遇到的坑及解决方法

  1. version_id_col 在 Core 层 update() 中不生效
  • 现象:用 session.execute(update(Product).where(...).values(inventory=...)) 时版本号不自动增加。
  • 根因:version_id_col 是 ORM 层的特性,只在 Session flush 对象时生效。Core 层的 update() 不经过 ORM 的 mapper。
  • 解决:要么用 ORM 对象修改属性(然后 flush),要么在 Core 层的 update 中手动加 values(version_id=Product.version_id + 1)
  1. PostgreSQL 的 REPEATABLE READ 不会产生 serialization failure
  • 现象:测试中 REPEATABLE READ 下并发更新同一行,第一个事务成功后第二个事务未报错而是静默覆盖。
  • 根因:PostgreSQL 的 REPEATABLE READ 使用"第一个提交者赢"策略,第二个提交者被静默阻塞直到第一个事务结束。如果两个事务同时 commit,第二个会被静默回滚。
  • 解决:加上 version_id_col 做乐观锁,不要仅依赖隔离级别。
  1. 事务中发送了外部 HTTP 请求导致长事务
  • 现象:下单事务里调了"发送微信通知",结果微信 API 超时 30 秒,事务全程持锁。
  • 解决:外部 IO 永远放在事务之外。先 commit 事务,再发通知。

测试验证

# tests/test_ch13_transaction.py
import pytest
from sqlalchemy import create_engine, String, Integer, Numeric, text, func, select
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, sessionmaker
from sqlalchemy.exc import StaleDataError

@pytest.fixture
def engine():
    return create_engine("sqlite:///:memory:", echo=False)

@pytest.fixture
def ProductModel(engine):
    class Base(DeclarativeBase):
        pass
    class Product(Base):
        __tablename__ = "products"
        id: Mapped[int] = mapped_column(primary_key=True)
        sku: Mapped[str] = mapped_column(String(30))
        inventory: Mapped[int] = mapped_column(Integer)
        version_id: Mapped[int] = mapped_column(Integer, server_default=text("1"))
        __mapper_args__ = {"version_id_col": version_id}
    Base.metadata.create_all(engine)
    return Product

def test_optimistic_lock_version_increment(engine, ProductModel):
    """验证版本号在每次更新时自动递增"""
    Factory = sessionmaker(bind=engine)
    with Factory() as s:
        p = ProductModel(sku="T1", inventory=10)
        s.add(p)
        s.commit()
        v1 = p.version_id
        p.inventory = 9
        s.commit()
        assert p.version_id == v1 + 1

def test_optimistic_lock_prevents_lost_update(engine, ProductModel):
    """验证版本号冲突时抛出 StaleDataError"""
    from sqlalchemy.orm import Session as OrmSession
    Factory = sessionmaker(bind=engine)
    with Factory() as s:
        p = ProductModel(sku="T2", inventory=5)
        s.add(p)
        s.commit()

    s1, s2 = Factory(), Factory()
    p1 = s1.get(ProductModel, p.id)
    p2 = s2.get(ProductModel, p.id)
    p1.inventory -= 1
    s1.commit()
    p2.inventory -= 1
    with pytest.raises(StaleDataError):
        s2.commit()
    s1.close()
    s2.rollback()
    s2.close()

def test_transaction_rollback(engine, ProductModel):
    """验证事务回滚后数据不变"""
    Factory = sessionmaker(bind=engine)
    with Factory() as s:
        p = ProductModel(sku="T3", inventory=100)
        s.add(p)
        s.commit()
        original = p.inventory
        p.inventory = 999
        s.rollback()
    with Factory() as s:
        reloaded = s.get(ProductModel, p.id)
        assert reloaded.inventory == original

四、项目总结

优点与缺点

对比维度 无事务管理 Session 事务 + 乐观锁
数据一致性 可能脏读、幻读、丢失更新 版本号 + 重试,保证数据一致
并发安全 超卖 StaleDataError 自动检测冲突
性能 无额外开销 版本号列略有开销,可忽略
死锁风险 乐观锁无死锁,悲观锁可能死锁

适用场景

  1. 库存扣减——必须保证不会超卖。
  2. 账户转账——余额增减必须原子。
  3. 座位预订/抢购——有限资源的高并发分配。
  4. 幂等下单——防止重复提交创建多余订单。
  5. 配置热更新——防止两人同时修改同一配置互相覆盖。

不推荐场景

  1. 纯日志/审计类追加写入——不存在并发冲突。
  2. 冲突概率极低的操作——引入版本号的收益小于复杂度。
  3. 用户级数据的独立修改——不同用户修改不同行不存在冲突。

注意事项

  1. version_id_col 每次 commit 都会递增:即使只是 select 后 commit(无实际变更),PostgreSQL 也有可能在 expire_on_commit 后触发 SELECT 而不是 UPDATE。观察实际 SQL。
  2. 重试次数和退避策略要合理:通常 3-5 次重试 + 指数退避(10ms, 100ms, 1s),超过则告警人工介入。
  3. 事务中不要做 I/O 密集操作:HTTP 调用、文件写入、消息队列发送都应放在事务结束后。
  4. 幂等键的设计:可以是客户端生成的 UUID,也可以是业务唯一键(如外部订单号)。必须有唯一约束保证数据库层也能防重。

常见踩坑经验

案例 1:乐观锁重试产生雪崩

  • 现象:秒杀时大量请求 StaleDataError,全部重试,数据库 CPU 100%。
  • 根因:10 个线程抢 1 件库存,每个线程都在密集重试——内耗大于实际有效请求。
  • 修复:引入抢购队列(Redis 排队),减小数据库的直接并发压力;或使用 SKIP LOCKED 跳过已锁行。

案例 2:幂等键设计不当导致重复创建

  • 现象:同一个幂等键在毫秒级间隔内被调了 2 次,仍然创建了 2 条订单。
  • 根因:幂等检查和插入之间没有原子性保证——两个并发请求都发现"不存在",都执行了 INSERT。
  • 修复:幂等键列加 UniqueConstraint,让第二个 INSERT 直接报 IntegrityError,业务层 catch 后返回已有订单。

案例 3:事务粒度太大导致连接池耗尽

  • 现象:下单接口里包裹了整个流程(查库存→创建订单→调用支付→发通知→生成物流单),事务耗时 3-5 秒,连接池瞬间耗尽。
  • 修复:把事务拆成多段——下单扣库存(1 个事务)→ 支付回调(另 1 个事务)→ 通知(事务外)。

思考题

  1. 请设计一个"秒杀"场景的下单方案。某商品库存仅有 100 件,预计 10000 人同时抢购。你会如何结合 Redis 预减库存 + 数据库乐观锁 + 消息队列来实现?画出架构图并解释每一层的职责。

  2. 如果你的系统使用 SERIALIZABLE 隔离级别搭配 PostgreSQL SSI,并发事务可能与乐观锁产生相同的冲突检测效果。那么什么时候仍然需要在 SERIALIZABLE 级别上再加版本号乐观锁?两者在冲突检测机制上的本质区别是什么?


参考答案参见附录 E。

延伸阅读与资源

NumPy 从入门到生产落地:全链路实战指南(科学计算/向量化)
Redis 8 实战精讲:从 CRUD 到源码,构建高可用缓存系统
Redis 实战修炼与原理进阶
Python 3实战精进:从脚本到高并发订单引擎
MongoDB 实战进阶与内核修炼
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析

Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐