目录
SeaORM 事务使用详解
SeaORM 主要提供两种事务写法:
transaction(...):闭包事务,闭包返回Ok自动提交,返回Err自动回滚。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()
})));
}
}
}
十三、常见错误
错误一:事务内混用 db 和 txn
错误:
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 泛型连接,确保普通连接和事务连接都能复用同一套数据库代码。