MongoDB: 事务处理:ACID与多文档一致性

最后更新:2026-08-26

事务保证多文档操作的原子性——掌握它能构建可靠的金融、订单系统。

1. 你将学到


100%
sequenceDiagram
    participant App as 应用
    participant DB as MongoDB<br/>副本集
    participant Log as 日志

    App->>DB: startTransaction()
    activate DB
    DB-->>App: session

    App->>DB: 扣款 $100
    DB-->>App: OK
    App->>DB: 加款 $100
    DB-->>App: OK
    App->>DB: 创建交易日志
    DB-->>App: OK

    alt 全部成功
        App->>DB: commitTransaction()
        DB-->>App: ✅ 提交
        DB->>Log: 持久化
    else 任何失败
        App->>DB: abortTransaction()
        DB-->>App: ❌ 回滚
        Note over DB: 所有变更撤销<br/>数据回到事务前
    end

    deactivate DB

2. 为什么需要事务?

概念说明:事务(Transaction)是数据库操作的逻辑单元,保证其中的所有操作要么全部成功、要么全部回滚。在涉及多个文档或多个集合的写入时(如转账、下单),没有事务就无法保证数据一致性——部分成功部分失败会导致脏数据。

工作原理:MongoDB 4.0+ 支持多文档 ACID 事务,基于 WiredTiger 引擎的快照隔离实现。事务开始时创建一个快照,所有读写操作基于该快照进行;提交时将变更原子性地应用到数据文件,回滚时丢弃所有变更。事务底层依赖副本集的 oplog 实现持久性和复制。

ACID 原理详解

特性 含义 MongoDB 实现 原理
Atomicity 原子性 事务全部成功或全部失败 commit / abort WiredTiger 通过 rollback log 保证:提交时原子写入,回滚时用 rollback log 恢复
Consistency 一致性 数据完整性约束不变 Schema Validation + 事务约束 事务提交前校验所有约束,违反则拒绝
Isolation 隔离性 并发事务互不干扰 快照隔离(Snapshot Isolation) 事务开始时获取数据快照,全程基于快照读写,不受其他事务影响
Durability 持久性 事务提交后永久保存 Journal + 副本集 oplog 写入先记 journal(WAL),再复制到副本集多数节点

MVCC 机制:MongoDB 通过多版本并发控制(MVCC)实现快照隔离。每个文档维护多个历史版本,事务根据开始时间戳读取对应版本的数据,写操作创建新版本而不覆盖旧版本,直到提交时新版本才对其他事务可见。

100%
graph TB
    subgraph "MVCC 多版本"
        D1["文档 v1<br/>balance: 1000"]
        D2["文档 v2<br/>balance: 900<br/>(事务A修改)"]
        D3["文档 v3<br/>balance: 1100<br/>(事务B修改)"]
    end

    subgraph "事务快照读取"
        T1["事务A (t1)<br/>读取 v1"] --> R1["balance: 1000"]
        T2["事务B (t2)<br/>读取 v1"] --> R2["balance: 1000"]
    end

    D1 --> D2
    D1 --> D3

    style D2 fill:#cce5ff
    style D3 fill:#d4edda

使用场景

JAVASCRIPT
// ❌ 反例:转账无事务
async function transfer(fromUserId, toUserId, amount) {
  await User.updateOne({ _id: fromUserId }, { $inc: { balance: -amount } });
  // 系统崩溃!
  await User.updateOne({ _id: toUserId }, { $inc: { balance: amount } });
  // 用户余额扣了但对方没收到
}

3. 事务基本用法

概念说明:MongoDB 事务通过 Session 对象管理——startSession() 创建会话,startTransaction() 开始事务,commitTransaction() 提交,abortTransaction() 回滚。所有事务内的操作必须传入 { session } 参数。

工作原理:事务的完整生命周期为:开始会话 → 开始事务 → 执行操作(带 session 参数)→ 提交/回滚 → 结束会话。提交时 WiredTiger 原子写入所有变更到 journal;回滚时通过 rollback log 撤销所有变更。事务超时默认 60 秒,超时自动回滚。

事务生命周期

100%
stateDiagram-v2
    [*] --> StartSession: startSession()
    StartSession --> Active: startTransaction()
    Active --> Active: 执行操作(带session)
    Active --> Committed: commitTransaction()
    Active --> Aborted: abortTransaction()
    Committed --> [*]: endSession()
    Aborted --> [*]: endSession()

    note right of Active: 默认60秒超时自动回滚
    note right of Committed: 变更持久化到journal

语法规则

步骤 方法 说明
1 startSession() 创建会话
2 startTransaction() 开始事务
3 操作 + {session} 所有读写必须传 session
4a commitTransaction() 全部成功 → 提交
4b abortTransaction() 任何失败 → 回滚
5 endSession() 释放会话资源
JAVASCRIPT
// === MongoDB 4.0+ 多文档事务 ===
const session = db.getMongo().startSession();

session.startTransaction();
try {
  // 1. 扣款
  db.users.updateOne(
    { _id: fromUserId },
    { $inc: { balance: -amount } },
    { session }
  );

  // 2. 加款
  db.users.updateOne(
    { _id: toUserId },
    { $inc: { balance: amount } },
    { session }
  );

  // 3. 提交事务
  await session.commitTransaction();
} catch (err) {
  // 4. 回滚
  await session.abortTransaction();
  throw err;
} finally {
  session.endSession();
}

要点解析

  1. 忘记传 { session } 的操作不在事务内,不受事务保护
  2. commitTransactionabortTransaction 是幂等操作,重复调用不会出错
  3. 事务超时后自动回滚,应用层应设置合理超时并做好重试

4. ACID 特性

概念说明:ACID 是数据库事务的四大核心保证——原子性(Atomicity)、一致性(Consistency)、隔离性(Isolation)、持久性(Durability)。理解 ACID 的 MongoDB 实现原理,是设计可靠事务系统的基础。

隔离级别详解:MongoDB 支持三种读隔离级别,通过 readConcern 控制:

隔离级别 readConcern 行为 使用场景
读未提交 local 读本地最新数据(可能回滚) 默认,性能优先
读已提交 majority 读已大多数确认的数据 强一致性需求
快照隔离 snapshot 事务内读一致的快照 事务内默认
100%
sequenceDiagram
    participant T1 as 事务1
    participant T2 as 事务2
    participant DB as MongoDB

    Note over DB: 初始 balance=1000

    T1->>DB: startTransaction(readConcern: snapshot)
    T1->>DB: 读 balance → 1000

    T2->>DB: startTransaction()
    T2->>DB: balance -100 → 写入 900
    T2->>DB: commitTransaction()

    T1->>DB: 读 balance → 1000(快照隔离,看不到T2的修改)

    Note over T1: 快照保证事务内一致性

    T1->>DB: commitTransaction()
    Note over DB: 冲突检测 → 如果T1也修改balance则报错
特性 含义 MongoDB 实现
Atomicity 原子性 事务全部成功或全部失败 commit / abort
Consistency 一致性 数据完整性约束 Schema Validation + 事务
Isolation 隔离性 并发事务互不干扰 快照隔离
Durability 持久性 事务提交后永久保存 Journal + 副本集

5. readConcern / writeConcern / readPreference

概念说明:这三个配置是 MongoDB 事务一致性调节的"三件套",分别控制读一致性、写持久性和读路由策略。它们决定了事务在一致性与性能之间的权衡点。

工作原理

(1) write Concern 写入关注

原理:写操作写入 Primary 后,需等待指定数量的 Secondary 确认复制才返回成功。

参数 行为 一致性 性能
w 1 仅 Primary 确认 最快
w majority 大多数节点确认 较慢
j true 写入磁盘 journal 最高 最慢
wtimeout ms 等待超时时间 超时返回错误
JAVASCRIPT
session.startTransaction({
  writeConcern: {
    w: 'majority',         // 大多数节点确认
    j: true,               // 写入磁盘 journal
    wtimeout: 5000         // 5 秒超时
  }
});

(2) read Concern 读取关注

原理:控制读操作能看到哪些版本的数据——是本地最新(可能未提交),还是多数确认的(已提交)。

level 说明 适用场景
local 读本地最新数据(默认) 性能优先,允许读到未提交数据
majority 读已大多数确认的数据 强一致性
snapshot 快照隔离(仅事务内) 事务默认,避免幻读
JAVASCRIPT
session.startTransaction({
  readConcern: {
    level: 'majority'      // 读已提交的数据
  },
  writeConcern: { w: 'majority' }
});

(3) read Preference 读取偏好

原理:控制读操作路由到 Primary 还是 Secondary,实现读写分离。

模式 行为 适用场景
primary 仅读主节点 强一致性事务
primaryPreferred 优先主,不可用则从 一般场景
secondary 仅读从节点 报表/分析,减轻主节点
secondaryPreferred 优先从,不可用则主 读多写少
nearest 网络延迟最低 地理分布式集群
JAVASCRIPT
session.startTransaction({
  readPreference: 'primary'             // 仅读主节点
});

session.startTransaction({
  readPreference: 'secondary'           // 读从节点
});

session.startTransaction({
  readPreference: 'secondaryPreferred'  // 优先从节点
});

三件套组合推荐

场景 writeConcern readConcern readPreference
金融事务 majority + j:true snapshot primary
一般事务 majority majority primary
报表分析 local secondary
开发测试 w:1 local primaryPreferred

6. mongoose 事务封装

概念说明:mongoose 提供了更优雅的事务 API——startSession() + session 参数。但手动管理事务的 try-catch-commit-abort 模板代码冗长,封装成 withTransaction 工具函数可大幅简化业务代码。

工作原理:mongoose 的 withTransaction 封装将 session 生命周期管理(start → commit/abort → end)自动化,业务函数只需关注核心逻辑。mongoose 还支持在 Model 操作(findByIdcreateupdateOne)中传入 { session } 参数,使事务内的 CRUD 无缝集成。

事务封装模式对比

模式 代码量 错误处理 重试支持 适用场景
手动 try-catch 手动 简单场景
withTransaction 封装 自动 可添加 生产推荐
mongoose.connection.transaction 最少 自动 内置 mongoose 6+
JAVASCRIPT
// === mongoose 事务封装 ===
async function withTransaction(callback) {
  const session = await mongoose.startSession();
  session.startTransaction();
  try {
    const result = await callback(session);
    await session.commitTransaction();
    return result;
  } catch (err) {
    await session.abortTransaction();
    throw err;
  } finally {
    session.endSession();
  }
}

// === 使用:转账 ===
async function transfer(fromUserId, toUserId, amount) {
  return withTransaction(async (session) => {
    const fromUser = await User.findById(fromUserId).session(session);
    if (fromUser.balance < amount) {
      throw new Error('Insufficient balance');
    }

    await User.updateOne(
      { _id: fromUserId },
      { $inc: { balance: -amount } },
      { session }
    );

    await User.updateOne(
      { _id: toUserId },
      { $inc: { balance: amount } },
      { session }
    );

    await TransactionLog.create([{
      fromUserId,
      toUserId,
      amount,
      createdAt: new Date()
    }], { session });

    return { success: true };
  });
}

▶ 示例 1:mongoose 事务重试封装

JAVASCRIPT
// Alice 的 ShopHub 金融系统:事务遇到 WriteConflict 自动重试
async function withRetryTransaction(callback, maxRetries = 3) {
  let lastError;
  for (let i = 0; i < maxRetries; i++) {
    const session = await mongoose.startSession();
    session.startTransaction({
      readConcern: { level: 'snapshot' },
      writeConcern: { w: 'majority' }
    });
    try {
      const result = await callback(session);
      await session.commitTransaction();
      return result;
    } catch (err) {
      await session.abortTransaction();
      lastError = err;
      if (err.errorLabels && err.errorLabels.includes('TransientTransactionError')) {
        console.log(`Retry ${i + 1}/${maxRetries} due to WriteConflict`);
        continue;
      }
      throw err;
    } finally {
      session.endSession();
    }
  }
  throw lastError;
}

// 使用
await withRetryTransaction(async (session) => {
  await User.updateOne({ _id: fromId }, { $inc: { balance: -100 } }, { session });
  await User.updateOne({ _id: toId }, { $inc: { balance: 100 } }, { session });
});

输出:

TEXT 📖 仅展示
// mongoose 操作成功执行
// 数据库查询/更新结果

7. 事务限制

概念说明:MongoDB 事务有明确的使用边界——必须副本集、有大小限制、不支持某些操作。理解这些限制是避免生产事故的关键。

限制详解

限制 说明 原因 应对策略
必须副本集 单机 MongoDB 不支持事务 事务依赖 oplog 实现持久性 开发环境可单节点副本集
16MB 文档 事务内的所有操作总和 WiredTiger 单文档限制 拆分大事务为小事务
默认 60 秒超时 maxTransactionLockRequestTimeoutMillis 避免长事务占用锁 调整超时参数
不能操作 capped collection 部分限制 capped 不支持回滚 避免事务内操作 capped
不能在事务内创建集合 部分限制(4.4+ 放宽) DDL 与事务冲突 事务前创建集合
写冲突 并发修改同一文档 乐观锁机制 自动重试 TransientTransactionError
锁等待 长事务阻塞其他操作 意图写锁 缩短事务,避免耗时操作

事务性能影响

操作 非事务 事务内 开销原因
单文档写入 基准 +30-50% 快照维护 + 锁管理
多文档写入 N次独立IO 1次提交 事务合并IO反而可能更快
读取 基准 +10-20% 快照读取额外开销
提交 5-50ms journal fsync + oplog 写入

事务最佳实践

实践 说明
事务尽量短 避免长事务占用锁,控制在 100ms 内
避免事务内计算 复杂计算放到事务外,事务内只做读写
重试 WriteConflict MongoDB 4.0+ 提供 errorLabels 标识可重试错误
优先用单文档原子操作 updateOne + $inc 本身原子,无需事务

8. Causal Consistency 因果一致性

概念说明:因果一致性(Causal Consistency)是比强一致性更轻量的一致性模型——它不保证所有操作全局有序,但保证有因果关系的操作按正确顺序执行。例如"先读余额,再扣款"——扣款操作必须基于最新读取的余额,这就是因果依赖。

工作原理:MongoDB 通过 operationTimeclusterTime 实现因果一致性。Session 内的每个操作会携带上次操作的 logicalTime,服务器确保后序操作看到前序操作的写入结果。启用因果一致性需要 readConcern: majority + writeConcern: majority

因果一致性 vs 其他一致性模型

模型 保证 性能 适用
强一致性(linearizable) 全局有序 最慢 金融核心
因果一致性 因果有序 较快 多步操作
最终一致性 无序 最快 日志、通知
读己之写 自己的写入可见 用户体验
100%
sequenceDiagram
    participant A as Alice
    participant P as Primary
    participant S as Secondary

    A->>P: 读余额 (readConcern: majority)
    P-->>A: balance=1000, clusterTime=t1

    A->>P: 扣款 $100 (writeConcern: majority)
    Note over A,P: 携带 afterClusterTime=t1
    P->>S: 复制 oplog
    S-->>P: 确认
    P-->>A: OK, clusterTime=t2

    A->>P: 查询交易记录 (readConcern: majority)
    Note over A,P: 携带 afterClusterTime=t2
    P-->>A: 包含刚扣款的记录 ✅

    Note over A,P: 因果一致:读一定能看到自己之前的写
JAVASCRIPT
// === 因果一致性:保证操作顺序 ===
const session = db.getMongo().startSession();
session.startTransaction({
  readConcern: { level: 'majority' },
  writeConcern: { w: 'majority' }
});

// 操作 1:读当前余额
const account = db.accounts.findOne({ userId: 'user_001' }, { session });
// 操作 2:基于读结果写入
db.accounts.updateOne(
  { userId: 'user_001' },
  { $set: { balance: account.balance - 100 } },
  { session }
);
// 保证:操作 2 看到的是操作 1 之后的状态

要点解析

  1. 因果一致性必须使用 Session,且 readConcernwriteConcern 都设为 majority
  2. 跨节点读取时(readPreference: secondary),因果一致性确保读到己方写入
  3. 因果一致性是 MongoDB 多文档事务和 Change Streams 的基础机制

▶ 示例 2:电商订单事务完整实战

JAVASCRIPT
// 场景:下单流程(订单 + 扣库存 + 钱包扣款 + 日志记录),全部原子性
// 前置:必须有副本集,事务才能工作

// 初始化数据
db.products.insertOne({ sku: 'PHONE-001', stock: 10, price: 599 });
db.users.insertOne({ _id: 'user_001', balance: 1000 });
db.transaction_logs.createIndex({ userId: 1, createdAt: -1 });

// 完整事务函数
async function placeOrder(userId, items) {
  const session = db.getMongo().startSession();
  session.startTransaction({
    readConcern: { level: 'snapshot' },
    writeConcern: { w: 'majority' }
  });

  try {
    // 1. 计算总金额 + 检查库存(原子读取)
    let total = 0;
    for (const item of items) {
      const product = db.products.findOne(
        { sku: item.sku, stock: { $gte: item.qty } },
        { session }
      );
      if (!product) {
        throw new Error(`库存不足: ${item.sku}`);
      }
      total += product.price * item.qty;
    }

    // 2. 检查用户余额
    const user = db.users.findOne({ _id: userId }, { session });
    if (user.balance < total) {
      throw new Error('余额不足');
    }

    // 3. 扣库存(带条件,防止超卖)
    for (const item of items) {
      const result = db.products.updateOne(
        { sku: item.sku, stock: { $gte: item.qty } },
        { $inc: { stock: -item.qty } },
        { session }
      );
      if (result.modifiedCount === 0) {
        throw new Error(`库存扣减失败: ${item.sku}`);
      }
    }

    // 4. 扣用户余额
    db.users.updateOne(
      { _id: userId, balance: { $gte: total } },
      { $inc: { balance: -total } },
      { session }
    );

    // 5. 创建订单
    const orderResult = db.orders.insertOne({
      userId,
      items,
      total,
      status: 'paid',
      createdAt: new Date()
    }, { session });

    // 6. 记录事务日志
    db.transaction_logs.insertOne({
      userId,
      orderId: orderResult.insertedId,
      amount: total,
      type: 'purchase',
      createdAt: new Date()
    }, { session });

    // 7. 提交事务
    session.commitTransaction();
    return { success: true, orderId: orderResult.insertedId };

  } catch (err) {
    // 任何失败 → 全部回滚
    session.abortTransaction();
    return { success: false, error: err.message };
  } finally {
    session.endSession();
  }
}

// 执行:下单
placeOrder('user_001', [
  { sku: 'PHONE-001', qty: 1 }
]);

// 测试回滚场景:故意制造错误
placeOrder('user_001', [
  { sku: 'NONEXIST', qty: 1 }  // 商品不存在
]);
// 抛出异常 → 事务回滚 → 库存、余额、订单、日志全部不变

// 验证原子性:
// db.products.findOne({ sku: 'PHONE-001' }) → stock: 10(未扣减)
// db.users.findOne({ _id: 'user_001' }) → balance: 1000(未扣款)

输出:事务成功时所有变更一起提交;事务失败时所有变更回滚,保证数据一致性。

▶ 示例 3:事务 + Change Stream 事件发布模式

事务保证数据一致性,但很多业务场景需要在事务提交后触发副作用——发送通知、更新缓存、写入审计日志。直接在事务内执行副作用很危险(事务回滚但通知已发出)。正确的模式是:事务内先写入"待发布事件"到 Event 集合,事务提交后用 Change Stream 监听 Event 集合,异步触发副作用。

JAVASCRIPT
// === 1. 事件集合 Schema ===
const eventSchema = new mongoose.Schema({
  type: { type: String, required: true },
  payload: mongoose.Schema.Types.Mixed,
  status: { type: String, enum: ['pending', 'published', 'failed'], default: 'pending' },
  createdAt: { type: Date, default: Date.now },
  publishedAt: { type: Date, default: null }
});
const Event = mongoose.model('Event', eventSchema);

// === 2. 事务内写入业务数据 + 事件 ===
async function transferBalance(fromId, toId, amount) {
  const session = await mongoose.startSession();
  session.startTransaction();
  try {
    const from = await User.findOneAndUpdate(
      { _id: fromId, balance: { $gte: amount } },
      { $inc: { balance: -amount } },
      { session, new: true }
    );
    if (!from) throw new Error('余额不足');

    const to = await User.findOneAndUpdate(
      { _id: toId },
      { $inc: { balance: amount } },
      { session, new: true }
    );

    // 事务内写入事件(事务提交后事件才可见)
    await Event.create([{
      type: 'balance.transferred',
      payload: { from: fromId, to: toId, amount, fromBalance: from.balance, toBalance: to.balance }
    }], { session });

    await session.commitTransaction();
    return { success: true, from, to };
  } catch (err) {
    await session.abortTransaction();
    return { success: false, error: err.message };
  } finally {
    session.endSession();
  }
}

// === 3. Change Stream 监听事件集合,异步触发副作用 ===
const changeStream = Event.watch([{ $match: { operationType: 'insert' } }]);

changeStream.on('change', async (change) => {
  const event = change.fullDocument;
  try {
    switch (event.type) {
      case 'balance.transferred':
        // 发送转账通知(邮件/短信/Push)
        await sendNotification(event.payload.to, `收到转账 ${event.payload.amount} 元`);
        await sendNotification(event.payload.from, `已转出 ${event.payload.amount} 元`);
        // 更新缓存
        await cacheDel(`user:balance:${event.payload.from}`);
        await cacheDel(`user:balance:${event.payload.to}`);
        break;
    }
    // 标记事件已发布
    await Event.updateOne({ _id: event._id }, { status: 'published', publishedAt: new Date() });
  } catch (err) {
    // 标记事件发布失败,后续可重试
    await Event.updateOne({ _id: event._id }, { status: 'failed' });
    console.error('事件发布失败:', event.type, err.message);
  }
});

// === 4. 失败事件重试机制 ===
async function retryFailedEvents() {
  const failed = await Event.find({ status: 'failed' })
    .sort({ createdAt: 1 })
    .limit(100)
    .lean();
  for (const event of failed) {
    try {
      await publishEvent(event);
      await Event.updateOne({ _id: event._id }, { status: 'published', publishedAt: new Date() });
    } catch (err) {
      console.error(`重试失败: ${event._id}`, err.message);
    }
  }
}
// 定时重试:每 5 分钟扫描一次
setInterval(retryFailedEvents, 5 * 60 * 1000);

输出:事务内同时写入业务数据和事件文档,事务提交后 Change Stream 自动捕获新事件,异步触发通知和缓存更新。失败事件自动重试,保证"至少一次"投递。

事务 + 事件模式的要点:1. 事件与业务数据在同一事务中——保证原子性,事务回滚则事件也回滚,不会发出错误通知;2. Change Stream 只在事务提交后才触发——避免了事务内读取到未提交数据的问题;3. 事件必须有状态追踪——pending → published/failed,failed 可重试;4. 幂等性——副作用必须幂等(重复执行结果相同),因为 Change Stream 可能重复投递;5. 事件顺序——同一文档的事件按 oplog 顺序投递,跨文档的事件不保证顺序。

❓ 常见问题

Q 事务能保证数据强一致性吗?
A 副本集 + writeConcern majority + readConcern majority + readPreference primary,可以。
Q 事务性能差多少?
A 相比非事务慢 30-50%。事务涉及锁和快照。
Q 单机 MongoDB 能用事务吗?
A 不能。必须副本集或分片集群。

📖 小节


📝 作业

  1. 基础题(⭐):用 mongosh 实现转账事务(含 try-catch 回滚)。
  2. 基础题(⭐):用 mongoose 封装 withTransaction 工具函数。
  3. 进阶题(⭐⭐):实现订单事务(下单 + 扣库存 + 创建订单 + 清空购物车,全部原子)。
  4. 进阶题(⭐⭐):测试事务失败回滚(故意抛错验证原子性)。
  5. 挑战题(⭐⭐⭐):完整电商事务系统(订单 + 库存 + 钱包 + 日志),支持分布式回滚。
Web-Tutorial.com

Web-Tutorial 技术团队

由多位开发者共同维护的编程教程平台。每篇教程由对应领域的开发者编写和审核,确保内容准确可靠。如发现任何问题,欢迎向我们反馈。

100%

🙏 帮我们做得更好

我们是刚上线的编程教程站,几个人的小团队,精力有限。页面虽经检查,难免还有疏漏——链接失效、排版错乱、内容有误、语言生硬……

如果您发现了,麻烦告诉我们,我们会在收到反馈后第一时间进行修复,再次感谢您的光临 🙏