MongoDB 在 4.0 版本之前不支持多文档事务,这让很多人望而却步。毕竟转账这种操作,扣钱和加钱必须一起成功或一起失败。4.0 之后 MongoDB 终于支持了多文档 ACID 事务,虽然性能和真正的关系型数据库比还是有差距,但绝大多数业务场景够用了。
我第一次用 MongoDB 事务是在一个积分系统里,用户兑换礼品需要同时扣减积分和生成订单。不用事务的话,极端情况下可能出现积分扣了但订单没生成,或者反过来。上了事务之后心里踏实多了。
基本用法
MongoDB 的事务是基于 Session 的,Go 里用 WithTransaction 最方便:
package transaction
import (
"context"
"errors"
"time"
"go.mongodb.org/mongo-driver/bson"
"go.mongodb.org/mongo-driver/bson/primitive"
"go.mongodb.org/mongo-driver/mongo"
"go.mongodb.org/mongo-driver/mongo/options"
"go.mongodb.org/mongo-driver/mongo/readconcern"
"go.mongodb.org/mongo-driver/mongo/readpref"
"go.mongodb.org/mongo-driver/mongo/writeconcern"
)
type TransactionManager struct {
client *mongo.Client
db *mongo.Database
}
funcNewTransactionManager(client *mongo.Client, dbName string) *TransactionManager {
return &TransactionManager{
client: client,
db: client.Database(dbName),
}
}
// 执行事务
func(tm *TransactionManager)ExecuteTransaction(fn func(sessCtx mongo.SessionContext)error) error {
session, err := tm.client.StartSession()
if err != nil {
return err
}
defer session.EndSession(context.Background())
_, err = session.WithTransaction(context.Background(), func(sessCtx mongo.SessionContext)(interface{}, error) {
returnnil, fn(sessCtx)
})
return err
}
转账是最经典的事务例子:
type Account struct {
ID primitive.ObjectID `bson:"_id,omitempty"`
UserID int64`bson:"user_id"`
Balance float64`bson:"balance"`
}
func(tm *TransactionManager)Transfer(fromUserID, toUserID int64, amount float64)error {
return tm.ExecuteTransaction(func(sessCtx mongo.SessionContext)error {
accountCol := tm.db.Collection("accounts")
// 1. 检查转出账户余额
var fromAccount Account
err := accountCol.FindOne(sessCtx, bson.M{"user_id": fromUserID}).Decode(&fromAccount)
if err != nil {
return err
}
if fromAccount.Balance < amount {
return errors.New("余额不足")
}
// 2. 扣款
_, err = accountCol.UpdateOne(
sessCtx,
bson.M{"user_id": fromUserID},
bson.M{"$inc": bson.M{"balance": -amount}},
)
if err != nil {
return err
}
// 3. 加款
_, err = accountCol.UpdateOne(
sessCtx,
bson.M{"user_id": toUserID},
bson.M{"$inc": bson.M{"balance": amount}},
)
if err != nil {
return err
}
// 4. 记录转账日志
_, err = tm.db.Collection("transfer_logs").InsertOne(sessCtx, bson.M{
"from_user_id": fromUserID,
"to_user_id": toUserID,
"amount": amount,
"created_at": time.Now(),
})
return err
})
}
WithTransaction 会自动处理重试,如果事务因为瞬态错误(比如主节点切换)失败,它会自动重试。但业务逻辑错误(比如余额不足)不会重试,直接返回错误。
手动控制事务
如果需要更精细的控制,可以手动开始、提交、回滚:
func(tm *TransactionManager)ManualTransaction()error {
ctx := context.Background()
session, err := tm.client.StartSession()
if err != nil {
return err
}
defer session.EndSession(ctx)
if err := session.StartTransaction(); err != nil {
return err
}
sessCtx := mongo.NewSessionContext(ctx, session)
accountCol := tm.db.Collection("accounts")
_, err = accountCol.UpdateOne(
sessCtx,
bson.M{"user_id": 1},
bson.M{"$inc": bson.M{"balance": -100}},
)
if err != nil {
session.AbortTransaction(sessCtx)
return err
}
_, err = accountCol.UpdateOne(
sessCtx,
bson.M{"user_id": 2},
bson.M{"$inc": bson.M{"balance": 100}},
)
if err != nil {
session.AbortTransaction(sessCtx)
return err
}
return session.CommitTransaction(sessCtx)
}
手动控制的好处是可以在事务中做更复杂的判断,比如在提交前再查一次确认状态。但大部分场景用 WithTransaction 就够了。
订单处理事务
电商下单是典型的事务场景,涉及库存、余额、订单三个 collection:
type Order struct {
ID primitive.ObjectID `bson:"_id,omitempty"`
OrderNo string`bson:"order_no"`
UserID int64`bson:"user_id"`
Items []OrderItem `bson:"items"`
TotalPrice float64`bson:"total_price"`
Status string`bson:"status"`
CreatedAt time.Time `bson:"created_at"`
}
type OrderItem struct {
ProductID int64`bson:"product_id"`
Quantity int`bson:"quantity"`
Price float64`bson:"price"`
}
func(tm *TransactionManager)CreateOrder(order *Order)error {
return tm.ExecuteTransaction(func(sessCtx mongo.SessionContext)error {
orderCol := tm.db.Collection("orders")
productCol := tm.db.Collection("products")
accountCol := tm.db.Collection("accounts")
// 1. 检查并扣减库存(乐观锁)
for _, item := range order.Items {
result, err := productCol.UpdateOne(
sessCtx,
bson.M{
"_id": item.ProductID,
"stock": bson.M{"$gte": item.Quantity},
},
bson.M{"$inc": bson.M{"stock": -item.Quantity}},
)
if err != nil {
return err
}
if result.ModifiedCount == 0 {
return errors.New("库存不足")
}
}
// 2. 扣减余额
result, err := accountCol.UpdateOne(
sessCtx,
bson.M{
"user_id": order.UserID,
"balance": bson.M{"$gte": order.TotalPrice},
},
bson.M{"$inc": bson.M{"balance": -order.TotalPrice}},
)
if err != nil {
return err
}
if result.ModifiedCount == 0 {
return errors.New("余额不足")
}
// 3. 创建订单
order.Status = "paid"
order.CreatedAt = time.Now()
_, err = orderCol.InsertOne(sessCtx, order)
return err
})
}
这里用了乐观锁思路:更新时带条件 stock >= quantity 和 balance >= totalPrice,如果条件不满足,ModifiedCount 就是 0,直接回滚事务。比先查再更新更安全,避免了并发下的竞态条件。
事务选项配置
默认的事务配置在大多数场景够用,但高并发或对一致性要求高的场景需要调整:
func(tm *TransactionManager)TransactionWithOptions(fn func(sessCtx mongo.SessionContext)error) error {
session, err := tm.client.StartSession()
if err != nil {
return err
}
defer session.EndSession(context.Background())
txnOpts := options.Transaction().
SetReadConcern(readconcern.Snapshot()).
SetWriteConcern(writeconcern.New(writeconcern.WMajority())).
SetReadPreference(readpref.Primary()).
SetMaxCommitTime(30 * time.Second)
_, err = session.WithTransaction(
context.Background(),
func(sessCtx mongo.SessionContext)(interface{}, error) {
returnnil, fn(sessCtx)
},
txnOpts,
)
return err
}
readconcern.Snapshot():读取事务开始时的快照,保证一致性writeconcern.WMajority():写入需要大多数节点确认,数据更安全MaxCommitTime:事务最大执行时间,防止长时间占用资源
读写关注级别根据业务需求选:
// 强一致性(金融场景)
rc := readconcern.Majority()
wc := writeconcern.New(
writeconcern.WMajority(),
writeconcern.J(true), // 等待日志刷盘
)
// 高性能(日志、监控)
rc2 := readconcern.Local()
wc2 := writeconcern.New(writeconcern.W(1))
重试与错误处理
MongoDB 事务可能会因为网络抖动、主节点切换等原因失败,需要实现重试机制:
func(tm *TransactionManager)ExecuteWithRetry(fn func(sessCtx mongo.SessionContext)error, maxRetriesint) error {
var err error
for i := 0; i < maxRetries; i++ {
err = tm.ExecuteTransaction(fn)
if err == nil {
returnnil
}
if !isRetryableError(err) {
return err
}
time.Sleep(time.Duration(i+1) * 100 * time.Millisecond)
}
return err
}
funcisRetryableError(err error)bool {
if err == nil {
returnfalse
}
if mongo.IsTransientTransactionError(err) {
returntrue
}
if mongo.IsUnknownTransactionCommitResult(err) {
returntrue
}
returnfalse
}
TransientTransactionError 表示事务执行过程中出错,可以重试;UnknownTransactionCommitResult 表示提交结果不确定,也可以重试。但如果是业务逻辑错误(如余额不足),重试多少次都没用。
带超时的执行:
func(tm *TransactionManager)ExecuteWithTimeout(fn func(sessCtx mongo.SessionContext)error, timeouttime.Duration) error {
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
session, err := tm.client.StartSession()
if err != nil {
return err
}
defer session.EndSession(ctx)
_, err = session.WithTransaction(ctx, func(sessCtx mongo.SessionContext)(interface{}, error) {
returnnil, fn(sessCtx)
})
return err
}
几个容易踩的坑
事务范围太大。MongoDB 事务有 60 秒默认超时,如果事务里做太多操作(比如批量更新几万条),很容易超时。建议把大事务拆成小事务,或者不用事务改用其他方案。
事务里查不到自己刚写入的数据。默认读关注级别下,事务内的写操作对事务内的读是可见的。但如果用了某些特殊的读关注,可能会出现读写不一致的情况。一般用 Snapshot 最稳妥。
事务和 bulkWrite 不能混用。事务里只能用单条操作(InsertOne、UpdateOne 等),不能用 BulkWrite。如果需要批量操作,在事务里循环执行单条操作。
分片集群的事务限制。MongoDB 4.2 之前分片集群不支持多文档事务,4.2 之后支持了但性能开销更大。如果用了分片,尽量减少跨片事务。
事务影响写入性能。开了事务之后,写入吞吐量会下降,因为需要记录 oplog、协调副本集。高并发写入场景要评估好性能影响。
忘记处理 TransientTransactionError。很多人写事务代码不判断错误类型,遇到可重试错误直接抛给上层。建议统一封装一个带重试的事务执行函数,业务代码只关心业务逻辑。
在事务外使用 session context。sessCtx 只能在 WithTransaction 的回调函数里用,拿出来用会报错。不要试图把 sessCtx 保存到外部变量里复用。
MongoDB 的事务虽然不如 PostgreSQL 那么成熟,但对于大多数业务场景已经够用了。关键是要控制事务范围,保持事务简短,避免长时间锁定文档。如果你发现事务经常超时或者性能不够,可能说明你的数据模型需要重新设计,比如把需要事务保证的数据放到同一个文档里,用 MongoDB 的原子单文档操作替代多文档事务。
夜雨聆风