Java 基础体系 · 第 89/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java MongoDB 工程:文档建模、驱动、事务、索引和变更流

MongoDB 是面向文档的数据库。它不是“把 JSON 放进数据库”,而是围绕 BSON 文档、集合、索引、复制集和分布式操作定义了一套一致性与查询模型。Java 工程要稳定使用 MongoDB,必须同时理解五个层次:

  1. 文档建模:决定数据如何组织,以及哪些不变量可以依靠单文档原子性维护。
  2. Java 驱动:负责连接池、BSON 编解码、会话、命令和错误传播。
  3. 事务:为跨文档、跨集合甚至跨分片的操作提供原子提交边界,但不能替代合理建模。
  4. 索引:改变查询路径和写入成本;索引不是“查询越多越好”。
  5. 变更流:把数据库中的变更以可恢复的事件流暴露给应用,但它不是消息队列,也不是天然的 exactly-once 处理系统。

本文以 Java 25 LTS 和 MongoDB Java Sync Driver 5.x 为范围。Java 驱动的具体 5.x 小版本应由项目统一依赖管理和安全策略决定;下面的 API 按 5.x 同步驱动接口书写。


一、MongoDB 的基本执行模型

1.1 文档、集合和 BSON

MongoDB 中的基本数据单元是 BSON 文档。BSON 是带有类型信息的二进制 JSON 风格格式,除了字符串、数字、数组和嵌套对象,还支持 ObjectId、日期、二进制数据和 Decimal128 等类型。

例如一个订单文档可以表示为:

{
  "_id": "order-1001",
  "customerId": "customer-7",
  "status": "PAID",
  "lines": [
    {
      "sku": "book-java",
      "quantity": 2,
      "unitPrice": 39.90
    }
  ],
  "totalAmount": 79.80,
  "createdAt": "2025-09-01T10:00:00Z"
}

集合是文档的容器。与关系数据库表相比,MongoDB 集合通常不要求所有文档具有同样的字段集合,但这不等于应用可以放任结构漂移。生产系统仍然需要:

  • Java 类型或 DTO 约束;
  • 必填字段和状态规则;
  • 数据库端 $jsonSchema 校验;
  • 迁移脚本和兼容读取逻辑。

MongoDB 单个 BSON 文档的大小上限是 16 MiB。因此,把无界增长的评论、日志、消息或历史状态全部嵌入一个文档,会最终触碰硬限制,而不是仅仅产生“文档变大”的性能问题。

1.2 单文档原子性是建模边界

MongoDB 对单个文档的写入具有原子性。假设订单文档内部有如下不变量:

totalAmount=i=1nquantityi×unitPricei\text{totalAmount}=\sum_{i=1}^{n} quantity_i \times unitPrice_i

如果订单行和订单总额都位于同一个文档中,那么修改订单行和总额可以由一次 updateOne 完成,其他读者不会看到只修改了一半的状态。

但如果订单和账户余额分别位于两个文档中:

orders/order-1001
accounts/customer-7

则一次普通的 updateOne 不能同时保证以下不变量:

account.balance=account.balanceorder.totalaccount.balance' = account.balance - order.total

此时有三个选择:

  1. 重新建模,把真正需要共同原子更新的数据放到一个文档;
  2. 接受最终一致性,并使用补偿或对账;
  3. 使用多文档事务。

事务是第三种选择,不是第一种选择的替代品。

1.3 MongoDB 与 JDBC、JPA 的边界

JDBC 面向关系数据库的表、行、列和 SQL;Jakarta Persistence(JPA)面向实体、关系映射、持久化上下文和 JPQL。MongoDB Java 驱动不实现 JDBC 的 ConnectionPreparedStatement 或 JPA 的 EntityManager 语义。

因此,以下概念不能直接类比:

关系数据库概念 MongoDB 中更接近的概念
集合
文档
字段
主键 通常是 _id
外键 通常由应用维护的引用字段
JOIN 聚合管道中的 $lookup,或应用侧查询
行锁 文档级原子更新、存储引擎内部并发控制
JDBC 事务 MongoDB session 上的多文档事务

这种差异会影响代码结构:MongoDB 应用通常直接使用 MongoDB 驱动、MongoTemplate 一类数据访问框架,或者其他明确支持 MongoDB 文档模型的库,而不是把关系型 DAO 强行套在文档上。


二、文档建模:嵌入、引用和不变量

2.1 嵌入模型

嵌入是把相关数据放在父文档中:

{
  "_id": "order-1001",
  "customerId": "customer-7",
  "lines": [
    { "sku": "book-java", "quantity": 2, "unitPrice": 39.90 },
    { "sku": "book-db", "quantity": 1, "unitPrice": 59.90 }
  ]
}

嵌入适用于以下条件:

  • 子数据通常和父数据一起读取;
  • 子数据数量有合理上界;
  • 子数据生命周期依附于父数据;
  • 需要在一次文档更新中维护它们的一致性。

嵌入后,读取一个订单不需要额外查询:

db.orders.findOne({ _id: "order-1001" })

更新一行时,可以使用数组过滤器:

db.orders.updateOne(
  { _id: "order-1001" },
  {
    $set: {
      "lines.$[line].quantity": 3
    }
  },
  {
    arrayFilters: [
      { "line.sku": "book-java" }
    ]
  }
)

这里的 $[line] 只匹配满足 arrayFilters 的数组元素。它不是把整个数组读回应用再写回,因此可以降低并发覆盖风险。

不过,嵌入并不意味着所有数组操作都便宜。对大型数组频繁使用 $push$pull 或数组内定位,会增加文档读写和索引维护成本。无界数组通常应该拆成独立集合,例如:

{
  "_id": "comment-9001",
  "postId": "post-1",
  "authorId": "user-2",
  "body": "...",
  "createdAt": "2025-09-01T10:00:00Z"
}

然后为 postId 和时间字段建立索引。

2.2 引用模型

引用模型只在父文档中保存关联标识:

{
  "_id": "order-1001",
  "customerId": "customer-7",
  "lineIds": ["line-1", "line-2"]
}

子文档位于其他集合。它适用于:

  • 子数据独立查询;
  • 子数据数量没有小而稳定的上界;
  • 子数据被多个父对象共享;
  • 子数据需要独立生命周期或权限;
  • 父文档接近 16 MiB 限制。

引用不是数据库外键。MongoDB 默认不会在删除客户时自动删除订单,也不会自动检查 customerId 是否存在。级联删除、孤儿检测和引用完整性必须由应用、事务、后台任务或数据库触发机制之外的流程维护。

2.3 读模式决定模型

关系建模经常从“实体之间是什么关系”出发;MongoDB 建模更应从“应用如何读取和更新”出发。

例如博客文章和评论有两种模型:

嵌入评论:

{
  "_id": "post-1",
  "title": "MongoDB",
  "comments": [
    { "author": "a", "body": "first", "createdAt": "..." }
  ]
}

当评论数量有限且文章详情页总是同时显示评论时,这种模型简单且具有单文档原子性。

独立评论集合:

{
  "_id": "comment-1",
  "postId": "post-1",
  "body": "first",
  "createdAt": "..."
}

当评论无限增长、需要分页、审核、独立搜索或按作者查询时,独立集合更合适。此时详情页查询文章后,再按:

db.comments.find({ postId: "post-1" })
  .sort({ createdAt: -1, _id: -1 })
  .limit(20)

获取评论。

分页使用 (createdAt, _id) 作为稳定游标,比 skip 更适合深页。若上一页最后一条记录为 (t, id),下一页条件可以写成:

{
  postId: "post-1",
  $or: [
    { createdAt: { $lt: t } },
    { createdAt: t, _id: { $lt: id } }
  ]
}

对应索引通常为:

{ postId: 1, createdAt: -1, _id: -1 }

2.4 反范式化的代价

反范式化是有意复制数据。例如订单中保存商品下单时的名称和价格:

{
  "sku": "book-java",
  "productNameSnapshot": "Java 工程实践",
  "unitPrice": 39.90
}

这是合理的,因为订单历史不应随着商品当前名称或价格变化而改变。复制的是“历史快照”,而不是无条件复制所有字段。

反例是把客户地址同时复制到订单、发票、配送和支付文档,却没有规定哪些字段是快照、哪些字段必须实时同步。结果会出现:

  • 新订单读取旧地址;
  • 更新客户地址时无法确定是否更新历史订单;
  • 多集合更新失败后产生部分复制;
  • 对账无法判断哪个值是权威值。

每个复制字段都应明确其语义:

currentAddress:当前权威地址
shippingAddressSnapshot:下单时地址快照

两者不能用同一个字段名混淆。


三、Java 25 中使用 MongoDB 驱动

3.1 依赖和客户端生命周期

使用 Maven 时,可以显式依赖同步驱动:

<dependency>
  <groupId>org.mongodb</groupId>
  <artifactId>mongodb-driver-sync</artifactId>
  <version>5.5.1</version>
</dependency>

5.5.1 只是示例版本。项目应根据 MongoDB Server 版本、Java 25 支持矩阵和组织的依赖审计策略选择经过验证的 5.x 版本,不应无条件复制示例版本。

MongoDB 客户端是线程安全的,通常应在应用进程内创建一个长期存活的 MongoClient

import com.mongodb.ConnectionString;
import com.mongodb.MongoClientSettings;
import com.mongodb.client.MongoClient;
import com.mongodb.client.MongoClients;

public final class MongoProvider implements AutoCloseable {
    private final MongoClient client;

    public MongoProvider(String uri) {
        var settings = MongoClientSettings.builder()
                .applyConnectionString(new ConnectionString(uri))
                .build();
        this.client = MongoClients.create(settings);
    }

    public MongoClient client() {
        return client;
    }

    @Override
    public void close() {
        client.close();
    }
}

正确的生命周期是:

进程启动
  -> 创建一个 MongoClient
  -> 多个请求共享 MongoClient
  -> 每个请求获取数据库和集合句柄
  -> 进程关闭时关闭 MongoClient

MongoDatabaseMongoCollection 是轻量级句柄,可以按需获取。不要每个 HTTP 请求创建和关闭 MongoClient,因为这会反复建立连接、执行服务器发现并制造连接池抖动。

同步驱动执行阻塞 I/O。Java 25 的虚拟线程可以降低大量阻塞任务占用的线程成本,但它不会扩大 MongoDB 服务端能力,也不会绕过驱动连接池上限。虚拟线程数量、连接池大小、数据库并发能力仍然需要一起控制。

3.2 URI、超时和连接池

一个可用于本地复制集的 URI 示例:

mongodb://app:secret@mongo1:27017,mongo2:27017/appdb?replicaSet=rs0&retryWrites=true&w=majority

关键参数的含义:

  • replicaSet=rs0:告诉驱动按复制集发现节点,事务和变更流通常依赖复制集拓扑;
  • retryWrites=true:对支持的写操作,在特定网络错误下自动重试;
  • w=majority:要求写入得到多数副本确认;
  • serverSelectionTimeoutMS:选择可用服务器的最长等待时间;
  • connectTimeoutMS:建立连接的超时;
  • socketTimeoutMS:套接字 I/O 超时。

超时不是越短越好。服务器选择超时过短会把瞬时选举误报为业务失败;套接字超时过长则可能让请求占用资源很久。应根据请求超时预算、负载均衡层超时和重试策略共同配置。

retryWrites=true 也不等于所有业务写入都可以安全重试。驱动只会对具有可重试语义的命令进行处理,应用仍必须处理唯一键冲突、事务提交结果未知和业务副作用重复等问题。

3.3 文档 API 与 POJO 编解码

直接使用 Document 适合动态结构或基础设施代码:

import com.mongodb.client.MongoCollection;
import org.bson.Document;

MongoCollection<Document> orders =
        database.getCollection("orders");

Document order = new Document("_id", "order-1001")
        .append("customerId", "customer-7")
        .append("status", "PAID")
        .append("lines", java.util.List.of(
                new Document("sku", "book-java")
                        .append("quantity", 2)
                        .append("unitPrice", 39.90)
        ));

orders.insertOne(order);

业务代码若到处使用字符串字段名,重命名字段、检查类型和表达状态会变得困难。POJO 编解码器可以把 Java 类型映射到 BSON:

import org.bson.codecs.pojo.annotations.BsonId;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.List;

public record Order(
        @BsonId String id,
        String customerId,
        String status,
        List<OrderLine> lines,
        BigDecimal totalAmount,
        Instant createdAt
) {}

public record OrderLine(
        String sku,
        int quantity,
        BigDecimal unitPrice
) {}

创建集合句柄时配置 POJO 编解码器:

import static org.bson.codecs.pojo.configuration.Conventions.DEFAULT_CONVENTIONS;

import org.bson.codecs.configuration.CodecRegistries;
import org.bson.codecs.pojo.PojoCodecProvider;

var pojoProvider = PojoCodecProvider.builder()
        .automatic(true)
        .conventions(DEFAULT_CONVENTIONS)
        .build();

var registry = CodecRegistries.fromRegistries(
        MongoClientSettings.getDefaultCodecRegistry(),
        CodecRegistries.fromProviders(pojoProvider)
);

MongoCollection<Order> orders =
        database.getCollection("orders", Order.class)
                 .withCodecRegistry(registry);

BigDecimalInstant 和自定义值对象必须验证驱动版本中的编码支持及其 BSON 表示。金额通常应使用 Decimal128 或最小货币单位整数,并禁止使用二进制浮点数 double 作为财务金额的权威值:

79.90 不是所有二进制浮点运算都能精确表示

如果金额使用整数分,则:

long totalCents = 2L * 3990L + 1L * 5990L;

若使用 Decimal128,则应定义舍入规则、货币精度和跨语言读取约定。数据库类型和 Java 类型必须形成明确契约。

3.4 查询、更新和结果检查

一个带条件的原子更新示例:

import static com.mongodb.client.model.Filters.*;
import static com.mongodb.client.model.Updates.*;

var result = orders.updateOne(
        and(eq("_id", "order-1001"), eq("status", "CREATED")),
        combine(
                set("status", "PAID"),
                currentDate("paidAt")
        )
);

if (result.getMatchedCount() == 0) {
    // 订单不存在,或状态已不是 CREATED
    throw new IllegalStateException("订单不能从当前状态支付");
}

这里的条件更新同时承担了乐观并发控制:

只有旧状态=CREATED 的文档才能转换为 PAID\text{只有旧状态}=CREATED\text{ 的文档才能转换为 }PAID

如果两个线程同时执行,最多一个线程能匹配 status: CREATED;另一个线程看到 matchedCount = 0。这比“先查询状态,再无条件更新”更可靠,因为后者在查询和更新之间存在竞态窗口。

对于计数器,应使用 $inc 而不是读改写:

collection.updateOne(
        eq("_id", "stock-book-java"),
        inc("available", -1)
);

但这只保证字段更新原子,不自动保证库存不能变负。可以将条件也放入更新过滤器:

var result = collection.updateOne(
        and(eq("_id", "stock-book-java"),
            gt("available", 0)),
        inc("available", -1)
);

if (result.getModifiedCount() != 1) {
    throw new IllegalStateException("库存不足");
}

四、事务:跨文档原子性及其限制

4.1 事务的前置条件

MongoDB 多文档事务通常需要:

  • 副本集或分片集群;
  • 支持事务的 MongoDB Server 版本;
  • 驱动通过 session 执行操作;
  • 事务中的读写使用同一个 session;
  • 事务的读偏好通常为 primary;
  • 合理的事务超时、写关注和读关注配置。

单机 standalone 部署不能直接等同于生产复制集环境。开发环境若要测试事务和变更流,应使用本地复制集,而不是只启动一个普通 mongod

事务解决的是以下数据流:

sequenceDiagram
    participant A as Java 应用
    participant S as ClientSession
    participant DB as MongoDB 复制集

    A->>S: startTransaction
    A->>DB: update orders (session)
    DB-->>S: 暂存写入
    A->>DB: update inventory (session)
    DB-->>S: 暂存写入
    A->>S: commitTransaction
    S->>DB: 提交
    DB-->>S: committed
    S-->>A: 成功

事务开始后,写入对事务外读者不可见,直到提交成功。提交失败时不能简单认为“事务一定没有执行”:网络可能在服务端提交完成后、客户端收到响应前中断,这就是“提交结果未知”。

4.2 Java 事务示例

以下代码把订单状态更新和库存扣减放在同一个事务中:

import com.mongodb.ReadConcern;
import com.mongodb.ReadPreference;
import com.mongodb.TransactionOptions;
import com.mongodb.WriteConcern;
import com.mongodb.client.ClientSession;
import org.bson.Document;

import static com.mongodb.client.model.Filters.*;
import static com.mongodb.client.model.Updates.*;

TransactionOptions options = TransactionOptions.builder()
        .readConcern(ReadConcern.SNAPSHOT)
        .writeConcern(WriteConcern.MAJORITY)
        .readPreference(ReadPreference.primary())
        .build();

try (ClientSession session = client.startSession()) {
    String orderId = "order-1001";
    String sku = "book-java";
    int quantity = 2;

    String result = session.withTransaction(() -> {
        var orderResult = orders.updateOne(
                session,
                and(eq("_id", orderId), eq("status", "CREATED")),
                combine(
                        set("status", "PAID"),
                        currentDate("paidAt")
                )
        );

        if (orderResult.getModifiedCount() != 1) {
            throw new IllegalStateException("订单不存在或状态不允许支付");
        }

        var stockResult = inventory.updateOne(
                session,
                and(eq("_id", sku), gte("available", quantity)),
                inc("available", -quantity)
        );

        if (stockResult.getModifiedCount() != 1) {
            throw new IllegalStateException("库存不足");
        }

        return "committed";
    }, options);

    System.out.println(result);
}

每一步成立的原因是:

  1. 订单过滤器要求状态为 CREATED,避免重复支付;
  2. 库存过滤器要求可用库存至少为 quantity,避免扣成负数;
  3. 两个操作共享同一个 ClientSession
  4. withTransaction 负责事务开始、正常提交以及符合驱动规范的事务重试路径;
  5. MAJORITY 使提交后的写入至少得到多数副本确认,具体可见性仍受读关注和读取节点影响。

业务异常应被抛出,让事务中止。不要捕获异常后返回“成功”,否则驱动会提交一个应用认为失败的事务。

4.3 事务回调必须可重入

事务辅助 API 可能在特定错误标签下重新执行事务回调。因此回调中的代码必须是数据库操作,而不应包含不可重复的外部副作用:

session.withTransaction(() -> {
    orders.insertOne(session, order);

    // 不安全:事务重试可能导致邮件重复发送
    mailClient.sendPaymentMail(order);

    return null;
}, options);

正确做法是把邮件请求写入事务内的 outbox 集合:

{
  "_id": "event-order-1001-paid",
  "type": "OrderPaid",
  "aggregateId": "order-1001",
  "payload": {},
  "processed": false
}

事务提交后,独立发送器读取 outbox 并发送邮件。发送器仍需要幂等键,因为“发送成功但处理标记未写回”同样会导致重试。

4.4 提交结果未知和重试边界

事务错误至少要区分:

  • TransientTransactionError:事务可能需要整体重试;
  • UnknownTransactionCommitResult:提交结果未知,通常应重试提交确认,而不是盲目重新执行所有业务副作用;
  • 唯一键冲突、校验失败、业务状态冲突:通常是确定性失败,不应无限重试。

具体重试策略应优先使用驱动提供的事务辅助 API,并设置有限次数、退避和总超时。应用日志应记录:

  • 事务业务标识;
  • session/事务相关上下文;
  • MongoDB 错误码和错误标签;
  • 是否已经进入提交阶段。

事务不能修复业务幂等性。比如客户端超时后重新提交“支付订单”请求,服务端可能已经完成第一次事务;订单状态条件或业务幂等键必须阻止第二次扣款。

4.5 事务的真实成本

事务会增加:

  • session 和事务状态管理;
  • 读写冲突重试;
  • 锁或冲突等待;
  • oplog 和复制压力;
  • 长事务占用资源的时间。

长事务会扩大冲突窗口,也可能阻碍历史数据清理或增加复制延迟。不要在事务中调用远程 HTTP、等待用户输入或执行耗时计算。事务内部只做必要的数据库读写,并使用请求级截止时间。


五、索引:从查询条件到执行计划

5.1 索引是什么

索引是按字段组织的辅助访问结构。没有合适索引时,MongoDB 可能扫描集合中的大量文档;有合适索引时,执行计划可以通过索引定位候选文档,再回表读取完整内容。

索引加速读的代价是:

  • 占用内存和磁盘;
  • 每次插入、删除和相关字段更新都要维护;
  • 建索引会消耗 CPU、I/O 和复制资源;
  • 多余索引会扩大写放大和启动恢复成本。

索引是否有用不能靠字段直觉判断,必须结合实际查询和 explain

5.2 唯一索引与业务约束

用户邮箱唯一性可以由唯一索引保证:

db.users.createIndex(
  { emailNormalized: 1 },
  { unique: true, name: "uq_users_email_normalized" }
)

应用应先规范化邮箱,例如统一大小写和空白,再写入 emailNormalized。只在应用层执行“先查询是否存在,再插入”不可靠,因为两个并发请求可能同时查询为空。唯一索引把约束下沉到数据库,第二个写入会收到 DuplicateKey 错误。

软删除场景可以使用部分唯一索引:

db.users.createIndex(
  { emailNormalized: 1 },
  {
    unique: true,
    partialFilterExpression: { deletedAt: { $exists: false } },
    name: "uq_active_users_email"
  }
)

这表示只有未删除文档参与唯一性约束。若历史数据中已有重复值,创建索引会失败;上线前必须先聚合检查重复记录并完成清理。

5.3 复合索引和前缀

假设查询是:

db.orders.find({
  customerId: "customer-7",
  status: "PAID"
}).sort({ createdAt: -1 }).limit(20)

可以建立:

db.orders.createIndex(
  { customerId: 1, status: 1, createdAt: -1 },
  { name: "ix_orders_customer_status_created" }
)

复合索引的前缀原则意味着该索引主要支持从左侧开始的字段组合,例如:

(customerId)
(customerId, status)
(customerId, status, createdAt)

它不等价于分别拥有 statuscreatedAt 的独立索引。字段顺序必须由真实过滤、排序和范围条件决定。

一个常用推导方式是:

  1. 先放选择性高且需要等值匹配的字段;
  2. 再考虑排序字段;
  3. 范围条件通常会限制后续字段对排序或定位的帮助;
  4. explain("executionStats") 验证实际扫描量。

例如:

db.orders.find({
  customerId: "customer-7",
  createdAt: { $gte: ISODate("2025-01-01T00:00:00Z") }
}).sort({ createdAt: -1 })
.explain("executionStats")

重点观察:

  • winningPlan 是否使用 IXSCAN
  • totalKeysExamined:扫描了多少索引键;
  • totalDocsExamined:读取了多少文档;
  • nReturned:最终返回多少文档。

如果返回 20 条却检查了数百万文档,索引可能没有覆盖主要过滤条件,或数据分布导致选择性很差。

5.4 多键索引和数组

数组字段上的索引称为多键索引。例如:

db.orders.createIndex(
  { "lines.sku": 1 },
  { name: "ix_orders_lines_sku" }
)

它会为数组中的元素建立索引入口。查询:

db.orders.find({ "lines.sku": "book-java" })

可以利用该索引。

但多个数组字段的组合会产生笛卡尔式索引条目风险,且某些复合多键索引受到“一个文档不能让多个索引字段同时来自数组”的结构限制。数组字段越多,越应该通过数据分布和 explain 验证,而不是凭字段名推断。

5.5 覆盖查询与投影

如果查询只需要索引中已有的字段,可以减少回表读取:

db.orders.find(
  { customerId: "customer-7", status: "PAID" },
  { _id: 1, createdAt: 1 }
)

_id 默认会返回;若索引不包含 _id,MongoDB 可能仍需读取文档,除非显式排除:

db.orders.find(
  { customerId: "customer-7", status: "PAID" },
  { _id: 0, createdAt: 1 }
)

是否形成真正的覆盖查询仍应由执行计划验证。投影本身不会自动让查询变快;如果查询阶段已经扫描大量文档,少返回几个字段不能解决索引问题。

5.6 索引上线和恢复

索引创建可能影响资源使用。生产变更应至少验证:

  1. 当前数据是否满足唯一索引约束;
  2. 创建过程的 CPU、磁盘、复制延迟;
  3. 创建完成后 listIndexes 是否包含预期索引;
  4. 目标查询的执行计划和延迟;
  5. 删除旧索引前是否仍有其他查询依赖它。

不要只根据索引名称判断它被使用。先从数据库分析、应用慢查询日志和 explain 找出查询形状,再决定创建或删除。


六、变更流:从 oplog 到可恢复事件

6.1 变更流是什么

变更流(Change Streams)允许应用订阅集合、数据库或整个部署范围内的变更。它建立在复制集或分片集群的 oplog 和服务器变更流机制之上。

典型事件路径是:

flowchart LR
    A[Java 应用写入] --> B[Primary]
    B --> C[Oplog]
    C --> D[Change Stream Cursor]
    D --> E[事件处理器]
    E --> F[下游索引/缓存/消息系统]
    E --> G[持久化 resume token]

变更流不是轮询集合。应用读取的是服务器生成的事件,并通过游标持续等待新事件。

常见事件类型包括:

  • insert
  • update
  • replace
  • delete
  • invalidate

drop、集合重建或部分拓扑变化可能导致游标失效。处理器不能只写一个“永远重连”的循环而不区分失效原因。

6.2 Java 订阅示例

import com.mongodb.client.ChangeStreamIterable;
import com.mongodb.client.model.Aggregates;
import com.mongodb.client.model.Filters;
import com.mongodb.client.model.changestream.FullDocument;
import com.mongodb.client.model.changestream.ChangeStreamDocument;
import org.bson.Document;

var pipeline = java.util.List.of(
        Aggregates.match(Filters.in(
                "operationType",
                java.util.List.of("insert", "update", "replace", "delete")
        ))
);

ChangeStreamIterable<Document> stream = orders.watch(pipeline)
        .fullDocument(FullDocument.UPDATE_LOOKUP);

try (var cursor = stream.cursor()) {
    while (cursor.hasNext()) {
        ChangeStreamDocument<Document> event = cursor.next();

        System.out.printf(
                "operation=%s key=%s fullDocument=%s%n",
                event.getOperationType(),
                event.getDocumentKey(),
                event.getFullDocument()
        );

        var token = event.getResumeToken();
        // 将 token 与下游处理状态按业务策略持久化
    }
}

fullDocument(UPDATE_LOOKUP) 的含义是:对更新事件,驱动请求服务器额外查询更新后的完整文档。它不保证这个文档仍然代表事件发生瞬间的状态。如果事件产生后文档又被修改,查到的完整文档可能已经包含后续修改。

因此:

  • 只需要知道发生了什么时,使用 updateDescription 更准确;
  • 需要当前文档快照时,可以使用 fullDocument,但要接受查询时序语义;
  • 不要把 fullDocument 当成事件时刻的不可变快照。

更新事件可能包含:

{
  "updateDescription": {
    "updatedFields": {
      "status": "PAID"
    },
    "removedFields": ["temporaryFlag"]
  }
}

如果下游需要可靠重建状态,应明确选择“按事件应用增量”还是“按主键重新读取当前状态”。

6.3 resume token 和断线恢复

每个变更事件带有 resume token。处理器应在成功完成下游处理后保存 token。发生网络中断后,可以:

var resumed = orders.watch(pipeline)
        .resumeAfter(savedResumeToken)
        .fullDocument(FullDocument.UPDATE_LOOKUP);

关键顺序是:

读取事件
  -> 执行下游副作用
  -> 确认副作用成功
  -> 持久化 resume token

如果先保存 token,再执行副作用,进程崩溃会跳过尚未处理的事件。如果先执行副作用,再保存 token,进程可能在两步之间崩溃,恢复后重复处理同一事件。

因此变更流通常只能自然地提供“至少一次”处理语义。要接近 exactly-once 的业务效果,需要下游幂等:

幂等键 = change event 的 resume token

或者使用业务事件 ID 建立唯一索引:

db.outbox_consumed.createIndex(
  { eventId: 1 },
  { unique: true }
)

下游处理和“标记已处理”最好放在同一个支持事务的存储边界内。若副作用发生在外部 HTTP 服务中,仍需要外部服务支持幂等键或采用 outbox/inbox 模式。

6.4 变更流的前置条件和丢失窗口

变更流依赖 oplog。若消费者停机时间超过 oplog 能覆盖的历史窗口,保存的 resume token 可能已经无法恢复,服务器会返回无法继续恢复的错误。此时处理器不能假装从断点继续,通常需要:

  1. 重新做一次全量同步;
  2. 记录新的同步边界;
  3. 再从该边界启动变更流;
  4. 通过幂等逻辑处理全量同步和增量事件之间的重叠。

如果数据库启用了集合级的变更前后镜像功能,MongoDB Server 版本和集合配置必须满足相应要求。例如创建集合时可以配置:

db.createCollection("orders", {
  changeStreamPreAndPostImages: { enabled: true }
})

这不是所有部署的默认能力,启用后还会增加存储和写入成本。Java 驱动能否读取相应字段取决于驱动和服务器版本,升级时应使用集成测试验证。


七、把订单、库存和索引组合成完整流程

下面构造一个有明确一致性边界的订单系统:

orders
  - 订单状态、订单行、金额快照

inventory
  - 每个 SKU 的可用库存

outbox
  - 已提交的业务事件

change stream
  - 将订单变化同步到搜索索引或缓存

7.1 集合结构

订单:

{
  "_id": "order-1001",
  "customerId": "customer-7",
  "status": "CREATED",
  "lines": [
    {
      "sku": "book-java",
      "quantity": 2,
      "unitPrice": 3990
    }
  ],
  "totalCents": 7980,
  "createdAt": ISODate("2025-09-01T10:00:00Z")
}

库存:

{
  "_id": "book-java",
  "available": 100,
  "version": 7
}

订单查询索引:

db.orders.createIndex(
  { customerId: 1, createdAt: -1, _id: -1 },
  { name: "ix_orders_customer_created" }
)

库存的 _id 已经是默认唯一索引,按 SKU 更新无需额外索引。

7.2 支付流程的数据变化

初始状态:

orders/order-1001.status = CREATED
inventory/book-java.available = 100

事务内执行:

1. 匹配 order-1001 且 status=CREATED
2. 将订单状态改为 PAID
3. 匹配 book-java 且 available >= 2
4. 将库存改为 98
5. 写入 outbox 事件
6. 提交事务

提交成功后:

orders/order-1001.status = PAID
inventory/book-java.available = 98
outbox/event-order-1001-paid 存在

如果第 3 步库存不足,事务回滚:

orders/order-1001.status 仍为 CREATED
inventory/book-java.available 不变
outbox 中没有已提交事件

如果订单状态更新成功但库存更新因过滤条件失败,而代码没有抛出异常,事务可能错误提交。因此,检查 matchedCountmodifiedCount 不是形式动作,而是业务状态转换的一部分。

7.3 为什么还需要 outbox

如果在订单事务提交后直接调用消息系统:

提交订单事务成功
  -> 调用消息系统失败

数据库已有支付订单,但下游没有收到事件。

如果先调用消息系统,再提交数据库:

消息发送成功
  -> 数据库事务失败

下游又会收到一个不存在的支付事件。

outbox 把业务状态和待发送事件放进同一个 MongoDB 事务,使二者一起提交。下游发送器可以轮询 outbox,或通过 outbox 集合的变更流触发发送。发送器要使用唯一事件 ID 和幂等处理,避免重启时重复产生不可逆副作用。


八、错误表现与诊断路径

8.1 “事务不工作”

常见表现:

MongoTransactionException
Transactions are not supported

诊断顺序:

  1. 确认连接到的是副本集或分片集群;
  2. 查看服务端版本;
  3. 确认所有事务操作都使用同一个 ClientSession
  4. 确认没有把操作发到 secondary;
  5. 查看错误码、错误标签和服务端日志;
  6. 检查事务是否超过超时或涉及不支持的操作组合。

本地测试可以启动复制集后连接:

mongodb://localhost:27017/?replicaSet=rs0

但仅在 URI 写上 replicaSet=rs0 不会把 standalone 自动变成复制集。

8.2 “明明有索引却很慢”

可能原因包括:

  • 查询字段顺序与复合索引前缀不匹配;
  • 过滤条件选择性很差;
  • 排序方向或字段组合不合适;
  • 数组多键索引产生大量候选;
  • 查询使用了隐式类型不匹配,例如字符串与整数;
  • 统计信息和实际数据分布使优化器选择了另一条计划;
  • 返回文档本身过大,瓶颈在网络和 BSON 解码,而不是索引定位。

诊断应保存真实查询的 explain("executionStats"),比较:

nReturned
totalKeysExamined
totalDocsExamined
executionTimeMillis

“使用了 IXSCAN”只说明访问了索引,不说明索引足够高效。

8.3 “变更流断线后重复或丢事件”

重复通常来自:

副作用成功
  -> 进程崩溃
  -> resume token 尚未保存
  -> 重启后再次消费

解决方向是幂等,而不是试图消灭所有崩溃窗口。

丢失通常来自:

  • 在副作用完成前保存 token;
  • resume token 超出 oplog 可恢复范围;
  • 收到 invalidate 后仍按普通网络错误重连;
  • 事件过滤和业务处理边界不一致。

处理器应记录:

最后读取 token
最后成功处理 token
最后持久化 token
下游副作用 ID

这些信息能区分“数据库没有发事件”“消费者没有读取”“读取后处理失败”和“处理成功但确认状态丢失”。

8.4 Java 类型和 BSON 类型不一致

例如历史数据中:

{ "totalCents": 7980 }
{ "totalCents": "7980" }

Java POJO 读取时可能出现解码异常,或者查询条件无法命中文档。修复方式不是在每个业务方法中添加隐式转换,而是:

  1. 统计字段类型分布;
  2. 编写迁移脚本统一 BSON 类型;
  3. 在写入入口强制类型;
  4. 增加数据库校验或集成测试。

MongoDB 的灵活 schema 是迁移自由度,不是取消数据契约。


九、生产边界与取舍

9.1 一致性选择必须显式化

写关注、读关注和读偏好共同决定观察到的行为:

  • WriteConcern.MAJORITY 主要约束写入确认;
  • ReadConcern.SNAPSHOT 用于事务快照读取;
  • ReadPreference.primary() 让读取面向主节点;
  • 从 secondary 读取可能存在复制延迟。

不能用“写入成功”直接推导“所有 secondary 都已经可读”,也不能用“事务提交成功”直接推导外部搜索索引已经更新。数据库事务边界和下游最终一致性边界是两件事。

9.2 不要把变更流当队列

变更流具有数据库变更语义:

  • 事件来源是 MongoDB 写入;
  • 生命周期受 oplog 和集合状态影响;
  • 事件消费进度由 resume token 表示;
  • 处理失败需要消费者自行重试和幂等。

它适合:

  • 缓存失效;
  • 搜索索引同步;
  • 审计投递;
  • 数据仓库增量采集;
  • 触发异步工作。

它不自动提供:

  • 任意长时间离线保存;
  • 消费者组负载均衡;
  • 消息延迟和重放策略;
  • 外部副作用 exactly-once。

需要这些能力时,应将变更流接入真正的消息系统,或使用 outbox 将业务事件写入明确的事件存储。

9.3 事务不是性能优化工具

如果每次更新订单都需要扫描多个集合、调用多个远程服务并等待很久,增加事务只会把问题集中到一个更大的失败单元。合理做法通常是:

单文档不变量 -> 单文档原子更新
有限跨文档不变量 -> 短事务
外部系统同步 -> outbox + 幂等消费者
历史和报表 -> 独立读模型或异步构建

这不是 MongoDB 专属规则,而是根据一致性边界划分故障边界。


十、最终检查:一个可验证的 MongoDB Java 工程应具备什么

一个可运行且可维护的 Java MongoDB 工程,至少应能回答以下问题:

  • 哪些字段共同组成一个单文档不变量?
  • 哪些数组有上界,哪些数组必须拆集合?
  • 引用关系的删除和完整性由谁维护?
  • Java 类型如何编码为 BSON,金额和时间的类型是什么?
  • MongoClient 是否是长期共享的?
  • 哪些操作需要 session,哪些操作可以单文档原子完成?
  • 事务回调是否可以安全重试?
  • 唯一性约束是否由唯一索引保证?
  • 每个高频查询使用什么索引,explain 的扫描比是多少?
  • 变更流的 resume token 保存在哪里?
  • 下游副作用如何幂等?
  • 消费者离线超过 oplog 窗口后如何全量重建?
  • 网络超时发生在提交前还是提交后,应用如何区分并恢复?

MongoDB 的工程能力不在于会调用 insertOne,而在于能把文档结构、原子性、驱动生命周期、索引执行计划和事件恢复语义组合成一个可验证的系统。 Java 25 提供现代的语言和并发运行环境,但数据库的一致性边界、重试边界和数据契约仍必须由应用明确设计。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。