不灭的焱

革命尚未成功,同志仍须努力 下载Java21

作者:AlbertWen  添加时间:2026-07-12 01:56:51  修改时间:2026-07-24 23:00:54  分类:01.Rust编程  编辑

目录

SeaORM 事务使用详解

SeaORM 主要提供两种事务写法:

  1. transaction(...):闭包事务,闭包返回 Ok 自动提交,返回 Err 自动回滚。
  2. begin():手动开启事务,再调用 commit()rollback()

SeaORM 官方更推荐大多数场景使用闭包事务;手动事务更适合复杂流程、提前回滚或闭包出现生命周期问题的情况。事务对象离开作用域而未提交时,会自动回滚。以下写法适用于 SeaORM 1.1.x 与当前 2.0.x 文档中的通用事务 API。(SeaQL)

一、事务解决什么问题

例如创建订单需要执行三步:

1. 创建订单
2. 创建订单明细
3. 扣减商品库存

假设:

  • 订单创建成功;
  • 订单明细创建成功;
  • 扣减库存失败。

如果没有事务,数据库里会留下一个不完整订单。

使用事务后,这三步被视为一个整体:

全部成功 → COMMIT
任一步失败 → ROLLBACK

二、准备工作

使用事务必须引入:

use sea_orm::TransactionTrait;

常见完整导入:

use sea_orm::{
    ActiveModelTrait,
    ColumnTrait,
    ConnectionTrait,
    DatabaseConnection,
    DatabaseTransaction,
    DbErr,
    EntityTrait,
    QueryFilter,
    QueryOrder,
    QuerySelect,
    Set,
    TransactionError,
    TransactionTrait,
};

以下案例假设项目中已经有这些 Entity:

entity/
├── account.rs
├── order.rs
├── order_item.rs
├── product.rs
└── transfer_log.rs

三、方式一:闭包事务 transaction()

1. 基本示例

创建用户时,同时创建用户资料。

use sea_orm::{
    ActiveModelTrait,
    DatabaseConnection,
    DbErr,
    Set,
    TransactionError,
    TransactionTrait,
};

use crate::entity::{user, user_profile};

pub async fn create_user_with_profile(
    db: &DatabaseConnection,
    username: String,
    display_name: String,
) -> Result<user::Model, TransactionError<DbErr>> {
    db.transaction::<_, user::Model, DbErr>(move |txn| {
        Box::pin(async move {
            // 1. 创建用户
            let user_model = user::ActiveModel {
                username: Set(username),
                status: Set(1),
                ..Default::default()
            }
            .insert(txn)
            .await?;

            // 2. 创建用户资料
            user_profile::ActiveModel {
                user_id: Set(user_model.id),
                display_name: Set(display_name),
                ..Default::default()
            }
            .insert(txn)
            .await?;

            // 返回 Ok:自动提交事务
            Ok(user_model)
        })
    })
    .await
}

闭包中的规则非常简单:

Ok(value)

代表:

COMMIT;

而:

Err(error)

代表:

ROLLBACK;

由于稳定版 Rust 的异步闭包支持仍有限,SeaORM 的闭包事务通常需要写成 Box::pin(async move { ... })。(SeaQL)

2. 泛型参数是什么意思

这段代码:

db.transaction::<_, user::Model, DbErr>(|txn| {
    // ...
})

三个泛型参数可以理解为:

transaction::<闭包类型, 成功返回类型, 业务错误类型>

也就是:

transaction::<_, user::Model, DbErr>

表示:

闭包类型:让编译器推断
成功返回:user::Model
失败返回:DbErr

如果不需要返回数据,可以写:

db.transaction::<_, (), DbErr>(|txn| {
    Box::pin(async move {
        // 执行数据库操作

        Ok(())
    })
})
.await?;

3. 最重要的规则:事务内必须使用 txn

正确写法:

user::ActiveModel {
    username: Set("albert".to_owned()),
    ..Default::default()
}
.insert(txn)
.await?;

错误写法:

user::ActiveModel {
    username: Set("albert".to_owned()),
    ..Default::default()
}
.insert(db) // 错误:使用了原始连接
.await?;

因为:

txn → 当前事务占用的数据库连接
db  → 连接池中的普通连接

如果事务内部误用了 db,这条 SQL 可能不属于当前事务,即使后面的事务回滚,这条数据也可能已经提交。

四、业务错误如何触发回滚

数据库错误可以直接使用 DbErr,但实际项目中经常还有业务错误,例如:

  • 库存不足;
  • 余额不足;
  • 用户状态异常;
  • 工单已经关闭;
  • 角色不能重复绑定。

推荐定义业务错误。

use sea_orm::DbErr;
use thiserror::Error;

#[derive(Debug, Error)]
pub enum OrderError {
    #[error("商品不存在:{0}")]
    ProductNotFound(i64),

    #[error("商品库存不足")]
    InsufficientStock,

    #[error("订单商品不能为空")]
    EmptyItems,

    #[error(transparent)]
    Database(#[from] DbErr),
}

注意:

#[error(transparent)]
Database(#[from] DbErr)

使数据库错误可以通过 ? 自动转换成 OrderError

创建订单事务示例

use sea_orm::{
    ActiveModelTrait,
    ColumnTrait,
    DatabaseConnection,
    EntityTrait,
    QueryFilter,
    Set,
    TransactionError,
    TransactionTrait,
};

use crate::entity::{order, order_item, product};

#[derive(Debug, Clone)]
pub struct CreateOrderItem {
    pub product_id: i64,
    pub quantity: i32,
}

pub async fn create_order(
    db: &DatabaseConnection,
    user_id: i64,
    items: Vec<CreateOrderItem>,
) -> Result<order::Model, TransactionError<OrderError>> {
    db.transaction::<_, order::Model, OrderError>(move |txn| {
        Box::pin(async move {
            if items.is_empty() {
                // 返回业务错误,整个事务自动回滚
                return Err(OrderError::EmptyItems);
            }

            // 1. 创建订单主表
            let order_model = order::ActiveModel {
                user_id: Set(user_id),
                status: Set("pending".to_owned()),
                ..Default::default()
            }
            .insert(txn)
            .await?;

            // 2. 处理订单明细
            for item in items {
                let product_model = product::Entity::find_by_id(item.product_id)
                    .one(txn)
                    .await?
                    .ok_or(OrderError::ProductNotFound(item.product_id))?;

                if product_model.stock < item.quantity {
                    return Err(OrderError::InsufficientStock);
                }

                // 3. 扣减库存
                let new_stock = product_model.stock - item.quantity;

                let mut product_active: product::ActiveModel = product_model.into();
                product_active.stock = Set(new_stock);
                product_active.update(txn).await?;

                // 4. 创建订单明细
                order_item::ActiveModel {
                    order_id: Set(order_model.id),
                    product_id: Set(item.product_id),
                    quantity: Set(item.quantity),
                    ..Default::default()
                }
                .insert(txn)
                .await?;
            }

            Ok(order_model)
        })
    })
    .await
}

任何位置出现:

return Err(...);

或者:

some_operation.await?;

返回错误,整个闭包事务都会回滚。

五、TransactionError 是什么

闭包事务的最终错误不是直接返回 OrderError,而是:

TransactionError<OrderError>

其主要包含两种错误:

pub enum TransactionError<E> {
    Connection(DbErr),
    Transaction(E),
}

含义如下:

Connection(DbErr)
    开启、提交、回滚事务时发生的数据库连接错误

Transaction(E)
    事务闭包内部返回的业务错误

这是 SeaORM 官方定义的事务错误结构。(Docs.rs)

调用时可以匹配:

match create_order(db, user_id, items).await {
    Ok(order) => {
        println!("订单创建成功:{}", order.id);
    }

    Err(TransactionError::Transaction(OrderError::EmptyItems)) => {
        println!("订单商品不能为空");
    }

    Err(TransactionError::Transaction(OrderError::InsufficientStock)) => {
        println!("库存不足");
    }

    Err(TransactionError::Transaction(err)) => {
        println!("订单业务执行失败:{err}");
    }

    Err(TransactionError::Connection(err)) => {
        println!("数据库事务连接失败:{err}");
    }
}

六、方式二:手动事务 begin()commit()rollback()

1. 最基本的手动事务

use sea_orm::{
    ActiveModelTrait,
    DatabaseConnection,
    DbErr,
    Set,
    TransactionTrait,
};

use crate::entity::{user, user_profile};

pub async fn create_user_manually(
    db: &DatabaseConnection,
) -> Result<(), DbErr> {
    // BEGIN
    let txn = db.begin().await?;

    user::ActiveModel {
        username: Set("albert".to_owned()),
        status: Set(1),
        ..Default::default()
    }
    .insert(&txn)
    .await?;

    user_profile::ActiveModel {
        user_id: Set(1001),
        display_name: Set("Albert".to_owned()),
        ..Default::default()
    }
    .insert(&txn)
    .await?;

    // COMMIT
    txn.commit().await?;

    Ok(())
}

注意手动事务通常传:

&txn

而闭包事务中,闭包参数本身已经是引用:

|txn| {
    // txn 类型是 &DatabaseTransaction
}

所以闭包事务中直接写:

.insert(txn)

2. 中途发生错误会怎么样

下面的代码中:

let txn = db.begin().await?;

first_operation(&txn).await?;
second_operation(&txn).await?;

txn.commit().await?;

假设:

second_operation(&txn).await?;

返回错误,函数会提前退出,txn 离开作用域。

SeaORM 会自动回滚没有提交的事务。(SeaQL)

因此这种写法是安全的:

pub async fn create_data(
    db: &DatabaseConnection,
) -> Result<(), DbErr> {
    let txn = db.begin().await?;

    create_first(&txn).await?;
    create_second(&txn).await?;

    txn.commit().await?;

    Ok(())
}

3. 显式调用 rollback()

需要在某个业务条件下立即终止时,可以显式回滚。

pub async fn deduct_stock(
    db: &DatabaseConnection,
    product_id: i64,
    quantity: i32,
) -> Result<(), OrderError> {
    let txn = db.begin().await?;

    let product_model = product::Entity::find_by_id(product_id)
        .one(&txn)
        .await?
        .ok_or(OrderError::ProductNotFound(product_id))?;

    if product_model.stock < quantity {
        txn.rollback().await?;

        return Err(OrderError::InsufficientStock);
    }

    let new_stock = product_model.stock - quantity;

    let mut product_active: product::ActiveModel = product_model.into();
    product_active.stock = Set(new_stock);
    product_active.update(&txn).await?;

    txn.commit().await?;

    Ok(())
}

对应 SQL 流程大致为:

BEGIN;

SELECT * FROM product WHERE id = ?;

-- 库存不足
ROLLBACK;

七、手动事务的推荐封装方式

业务代码比较复杂时,可以把事务内部逻辑拆成单独函数。

use sea_orm::DatabaseTransaction;

async fn create_order_inner(
    txn: &DatabaseTransaction,
    user_id: i64,
) -> Result<order::Model, OrderError> {
    let order_model = order::ActiveModel {
        user_id: Set(user_id),
        status: Set("pending".to_owned()),
        ..Default::default()
    }
    .insert(txn)
    .await?;

    // 其他事务操作……

    Ok(order_model)
}

外层控制事务:

pub async fn create_order_manual(
    db: &DatabaseConnection,
    user_id: i64,
) -> Result<order::Model, OrderError> {
    let txn = db.begin().await?;

    let order_model = create_order_inner(&txn, user_id).await?;

    txn.commit().await?;

    Ok(order_model)
}

职责非常清晰:

外层函数:开启、提交事务
inner 函数:处理业务逻辑

八、让 Service 同时支持普通连接和事务连接

这是 SeaORM 项目中非常实用的设计。

因为以下对象都实现了 ConnectionTrait

DatabaseConnection
DatabaseTransaction

所以 Service 函数不要固定接收:

&DatabaseConnection

而可以接收泛型连接:

&C
where
    C: ConnectionTrait

SeaORM 的查询和写入 API 普遍接受实现了 ConnectionTrait 的连接,因此同一个函数可以运行在普通连接或事务连接上。(Docs.rs)

示例

use sea_orm::{
    ActiveModelTrait,
    ConnectionTrait,
    DbErr,
    Set,
};

use crate::entity::operation_log;

pub async fn write_operation_log<C>(
    conn: &C,
    user_id: i64,
    action: &str,
) -> Result<(), DbErr>
where
    C: ConnectionTrait,
{
    operation_log::ActiveModel {
        user_id: Set(user_id),
        action: Set(action.to_owned()),
        ..Default::default()
    }
    .insert(conn)
    .await?;

    Ok(())
}

普通调用:

write_operation_log(db, user_id, "用户登录").await?;

事务中调用:

let txn = db.begin().await?;

write_operation_log(&txn, user_id, "创建订单").await?;

txn.commit().await?;

闭包事务中调用:

db.transaction::<_, (), DbErr>(|txn| {
    Box::pin(async move {
        write_operation_log(txn, user_id, "创建订单").await?;

        Ok(())
    })
})
.await?;

这种写法非常适合你的多模块项目:

modules/
├── user/
│   ├── service.rs
│   └── repository.rs
├── order/
│   ├── service.rs
│   └── repository.rs
└── system/
    └── operation_log_service.rs

九、转账事务完整示例

转账是最典型的事务场景:

1. 锁定转出账户
2. 锁定转入账户
3. 检查余额
4. 扣减转出账户余额
5. 增加转入账户余额
6. 记录流水
7. 提交事务

除了事务,还必须考虑并发。

假设两个请求同时读取余额:

账户余额:100

请求 A 读取到 100
请求 B 读取到 100

A 扣减 80
B 扣减 80

如果只使用普通查询,就可能出现并发覆盖。因此要使用:

SELECT ... FOR UPDATE

SeaORM 的 Select 查询提供 lock(LockType::Update) 等行锁 API。(Docs.rs)

1. 定义错误

use rust_decimal::Decimal;
use sea_orm::DbErr;
use thiserror::Error;

#[derive(Debug, Error)]
pub enum TransferError {
    #[error("转出账户和转入账户不能相同")]
    SameAccount,

    #[error("转账金额必须大于 0")]
    InvalidAmount,

    #[error("账户不存在:{0}")]
    AccountNotFound(i64),

    #[error("余额不足,当前余额:{balance},转账金额:{amount}")]
    InsufficientBalance {
        balance: Decimal,
        amount: Decimal,
    },

    #[error(transparent)]
    Database(#[from] DbErr),
}

2. 完整转账代码

use rust_decimal::Decimal;
use sea_orm::{
    sea_query::LockType,
    ActiveModelTrait,
    DatabaseConnection,
    EntityTrait,
    QuerySelect,
    Set,
    TransactionError,
    TransactionTrait,
};

use crate::entity::{account, transfer_log};

pub async fn transfer(
    db: &DatabaseConnection,
    from_account_id: i64,
    to_account_id: i64,
    amount: Decimal,
) -> Result<(), TransactionError<TransferError>> {
    db.transaction::<_, (), TransferError>(move |txn| {
        Box::pin(async move {
            if from_account_id == to_account_id {
                return Err(TransferError::SameAccount);
            }

            if amount <= Decimal::ZERO {
                return Err(TransferError::InvalidAmount);
            }

            /*
             * 按固定顺序锁账户,降低两个转账事务互相等待造成死锁的概率。
             *
             * 例如:
             * A -> B
             * B -> A
             *
             * 两个事务都先锁 ID 较小的账户。
             */
            let (first_id, second_id) =
                if from_account_id < to_account_id {
                    (from_account_id, to_account_id)
                } else {
                    (to_account_id, from_account_id)
                };

            let first_account = account::Entity::find_by_id(first_id)
                .lock(LockType::Update)
                .one(txn)
                .await?
                .ok_or(TransferError::AccountNotFound(first_id))?;

            let second_account = account::Entity::find_by_id(second_id)
                .lock(LockType::Update)
                .one(txn)
                .await?
                .ok_or(TransferError::AccountNotFound(second_id))?;

            // 恢复成转出账户和转入账户
            let (from_account, to_account) =
                if from_account_id == first_id {
                    (first_account, second_account)
                } else {
                    (second_account, first_account)
                };

            if from_account.balance < amount {
                return Err(TransferError::InsufficientBalance {
                    balance: from_account.balance,
                    amount,
                });
            }

            let new_from_balance = from_account.balance - amount;
            let new_to_balance = to_account.balance + amount;

            // 扣减转出账户余额
            let mut from_active: account::ActiveModel =
                from_account.into();

            from_active.balance = Set(new_from_balance);
            from_active.update(txn).await?;

            // 增加转入账户余额
            let mut to_active: account::ActiveModel =
                to_account.into();

            to_active.balance = Set(new_to_balance);
            to_active.update(txn).await?;

            // 记录转账流水
            transfer_log::ActiveModel {
                from_account_id: Set(from_account_id),
                to_account_id: Set(to_account_id),
                amount: Set(amount),
                status: Set("success".to_owned()),
                ..Default::default()
            }
            .insert(txn)
            .await?;

            Ok(())
        })
    })
    .await
}

执行逻辑相当于:

BEGIN;

SELECT *
FROM account
WHERE id = ?
FOR UPDATE;

SELECT *
FROM account
WHERE id = ?
FOR UPDATE;

UPDATE account
SET balance = balance - ?
WHERE id = ?;

UPDATE account
SET balance = balance + ?
WHERE id = ?;

INSERT INTO transfer_log (...);

COMMIT;

任何一步失败:

ROLLBACK;

十、嵌套事务

SeaORM 支持在事务中再次开启事务:

let outer_txn = db.begin().await?;
let inner_txn = outer_txn.begin().await?;

嵌套事务底层通常使用数据库的 SAVEPOINT 实现,而不是再开启一个完全独立的数据库事务。(SeaQL)

大致对应:

BEGIN;

SAVEPOINT savepoint_1;

-- 内层操作

ROLLBACK TO SAVEPOINT savepoint_1;

COMMIT;

嵌套事务示例

pub async fn nested_transaction_example(
    db: &DatabaseConnection,
) -> Result<(), DbErr> {
    // 外层事务
    let outer_txn = db.begin().await?;

    user::ActiveModel {
        username: Set("user-a".to_owned()),
        ..Default::default()
    }
    .insert(&outer_txn)
    .await?;

    {
        // 内层事务,对应 SAVEPOINT
        let inner_txn = outer_txn.begin().await?;

        user::ActiveModel {
            username: Set("user-b".to_owned()),
            ..Default::default()
        }
        .insert(&inner_txn)
        .await?;

        // 只回滚内层事务
        inner_txn.rollback().await?;
    }

    // user-a 会提交
    // user-b 已经回滚
    outer_txn.commit().await?;

    Ok(())
}

最终结果:

user-a:存在
user-b:不存在

内层提交后,外层还能回滚吗

可以。

let outer_txn = db.begin().await?;

let inner_txn = outer_txn.begin().await?;

create_data(&inner_txn).await?;

inner_txn.commit().await?;

// 外层回滚
outer_txn.rollback().await?;

即使内层调用了:

inner_txn.commit()

也只是释放对应的保存点。

如果外层最终回滚,内层已经执行的数据库修改仍然会被一起回滚。

十一、事务隔离级别

SeaORM 提供:

begin_with_config()
transaction_with_config()

可以指定:

IsolationLevel
AccessMode

当前官方文档说明,这部分配置主要针对 MySQL 和 PostgreSQL 实现。(SeaQL)

手动事务设置隔离级别

use sea_orm::{
    AccessMode,
    IsolationLevel,
    TransactionTrait,
};

pub async fn serializable_operation(
    db: &DatabaseConnection,
) -> Result<(), DbErr> {
    let txn = db
        .begin_with_config(
            Some(IsolationLevel::Serializable),
            Some(AccessMode::ReadWrite),
        )
        .await?;

    // 执行业务操作
    // ...

    txn.commit().await?;

    Ok(())
}

闭包事务设置隔离级别

use sea_orm::{
    AccessMode,
    DbErr,
    IsolationLevel,
    TransactionTrait,
};

pub async fn execute_with_config(
    db: &DatabaseConnection,
) -> Result<(), TransactionError<DbErr>> {
    db.transaction_with_config::<_, (), DbErr>(
        |txn| {
            Box::pin(async move {
                // 数据库操作
                // ...

                Ok(())
            })
        },
        Some(IsolationLevel::Serializable),
        Some(AccessMode::ReadWrite),
    )
    .await
}

常见隔离级别包括:

IsolationLevel::ReadUncommitted
IsolationLevel::ReadCommitted
IsolationLevel::RepeatableRead
IsolationLevel::Serializable

访问模式包括:

AccessMode::ReadOnly
AccessMode::ReadWrite

这些枚举及配置方法由 SeaORM 的 TransactionTrait 提供。(Docs.rs)

十二、在 Salvo 项目中的推荐结构

结合你之前的技术栈:

Rust 2024
Salvo
SeaORM
MySQL 8.0
Redis 6

建议让 Handler 不直接处理事务。

HTTP Handler
    ↓
Service:控制事务和业务规则
    ↓
Repository:执行数据库操作
    ↓
SeaORM Entity

目录示例:

src/
├── applications/
│   └── api/
│       └── handlers/
│           └── order_handler.rs
│
├── modules/
│   └── order/
│       ├── service.rs
│       ├── repository.rs
│       ├── dto.rs
│       └── error.rs
│
├── entities/
│   ├── order.rs
│   ├── order_item.rs
│   └── product.rs
│
└── infrastructure/
    └── database.rs

Repository 层

use sea_orm::{
    ActiveModelTrait,
    ConnectionTrait,
    DbErr,
    Set,
};

pub struct OrderRepository;

impl OrderRepository {
    pub async fn insert<C>(
        conn: &C,
        user_id: i64,
    ) -> Result<order::Model, DbErr>
    where
        C: ConnectionTrait,
    {
        order::ActiveModel {
            user_id: Set(user_id),
            status: Set("pending".to_owned()),
            ..Default::default()
        }
        .insert(conn)
        .await
    }
}

Service 层控制事务

use sea_orm::{
    DatabaseConnection,
    TransactionError,
    TransactionTrait,
};

pub struct OrderService;

impl OrderService {
    pub async fn create_order(
        db: &DatabaseConnection,
        user_id: i64,
    ) -> Result<order::Model, TransactionError<OrderError>> {
        db.transaction::<_, order::Model, OrderError>(move |txn| {
            Box::pin(async move {
                let order =
                    OrderRepository::insert(txn, user_id).await?;

                write_operation_log(
                    txn,
                    user_id,
                    "创建订单",
                )
                .await?;

                Ok(order)
            })
        })
        .await
    }
}

Salvo Handler

use salvo::prelude::*;

#[handler]
pub async fn create_order_handler(
    depot: &mut Depot,
    res: &mut Response,
) {
    let db = depot
        .obtain::<DatabaseConnection>()
        .expect("DatabaseConnection 未注入");

    let user_id = 1001;

    match OrderService::create_order(db, user_id).await {
        Ok(order) => {
            res.render(Json(serde_json::json!({
                "code": 0,
                "message": "订单创建成功",
                "data": {
                    "order_id": order.id
                }
            })));
        }

        Err(err) => {
            res.status_code(StatusCode::INTERNAL_SERVER_ERROR);

            res.render(Json(serde_json::json!({
                "code": 500,
                "message": err.to_string()
            })));
        }
    }
}

十三、常见错误

错误一:事务内混用 dbtxn

错误:

db.transaction::<_, (), DbErr>(|txn| {
    Box::pin(async move {
        create_order(txn).await?;

        // 不属于当前事务
        create_log(db).await?;

        Ok(())
    })
})
.await?;

正确:

db.transaction::<_, (), DbErr>(|txn| {
    Box::pin(async move {
        create_order(txn).await?;
        create_log(txn).await?;

        Ok(())
    })
})
.await?;

错误二:捕获错误后仍然返回 Ok

错误:

db.transaction::<_, (), DbErr>(|txn| {
    Box::pin(async move {
        if let Err(err) = create_order(txn).await {
            tracing::error!("创建订单失败:{err}");

            // 错误被吃掉,事务仍然会提交
        }

        Ok(())
    })
})
.await?;

正确:

db.transaction::<_, (), DbErr>(|txn| {
    Box::pin(async move {
        if let Err(err) = create_order(txn).await {
            tracing::error!("创建订单失败:{err}");

            // 把错误继续返回,事务才能回滚
            return Err(err);
        }

        Ok(())
    })
})
.await?;

或者直接:

create_order(txn).await?;

错误三:事务中执行耗时的外部操作

不建议:

let txn = db.begin().await?;

update_database(&txn).await?;

// 调用外部邮件、HTTP、AI 或文件服务,可能耗时几十秒
send_email().await?;
call_remote_api().await?;

txn.commit().await?;

事务持续期间可能一直占用数据库连接和行锁。

更合理的流程:

1. 在事务中写入业务数据
2. 在事务中写入待发送事件表
3. 提交事务
4. 后台任务发送邮件或调用外部 API

即 Outbox 思路:

db.transaction::<_, (), DbErr>(|txn| {
    Box::pin(async move {
        create_order(txn).await?;

        create_outbox_event(
            txn,
            "order.created",
        )
        .await?;

        Ok(())
    })
})
.await?;

错误四:用事务代替并发控制

事务本身不代表一定不会超卖。

错误逻辑:

let product = product::Entity::find_by_id(product_id)
    .one(txn)
    .await?;

if product.stock >= quantity {
    // 更新库存
}

多个事务可能同时读取到相同库存。

扣库存、转账、抢占资源等场景通常还需要:

.lock(LockType::Update)

或者使用带条件的原子更新:

UPDATE product
SET stock = stock - ?
WHERE id = ?
  AND stock >= ?;

然后检查影响行数。

十四、实际项目如何选择

普通的多表写入,优先使用闭包事务:

db.transaction(|txn| {
    Box::pin(async move {
        // ...
        Ok(())
    })
})
.await

流程复杂、需要分支控制或遇到生命周期问题时,使用手动事务:

let txn = db.begin().await?;

// ...

txn.commit().await?;

需要局部回滚时,使用嵌套事务:

let inner = outer.begin().await?;

转账、扣库存等高并发场景:

事务
+ 行锁
+ 固定加锁顺序
+ 数据库约束

最推荐的工程化方式是:Service 层控制事务,Repository 方法接收 ConnectionTrait 泛型连接,确保普通连接和事务连接都能复用同一套数据库代码。