乐于分享
好东西不私藏

MongoDB事务处理:多文档ACID保证

MongoDB事务处理:多文档ACID保证

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)errorerror {
 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)errorerror {
 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)errormaxRetriesinterror {
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)errortimeouttime.Durationerror {
 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 contextsessCtx 只能在 WithTransaction 的回调函数里用,拿出来用会报错。不要试图把 sessCtx 保存到外部变量里复用。


MongoDB 的事务虽然不如 PostgreSQL 那么成熟,但对于大多数业务场景已经够用了。关键是要控制事务范围,保持事务简短,避免长时间锁定文档。如果你发现事务经常超时或者性能不够,可能说明你的数据模型需要重新设计,比如把需要事务保证的数据放到同一个文档里,用 MongoDB 的原子单文档操作替代多文档事务。