Java 基础体系 · 第 89/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java MongoDB 工程:文档建模、驱动、事务、索引和变更流
MongoDB 是面向文档的数据库。它不是“把 JSON 放进数据库”,而是围绕 BSON 文档、集合、索引、复制集和分布式操作定义了一套一致性与查询模型。Java 工程要稳定使用 MongoDB,必须同时理解五个层次:
- 文档建模:决定数据如何组织,以及哪些不变量可以依靠单文档原子性维护。
- Java 驱动:负责连接池、BSON 编解码、会话、命令和错误传播。
- 事务:为跨文档、跨集合甚至跨分片的操作提供原子提交边界,但不能替代合理建模。
- 索引:改变查询路径和写入成本;索引不是“查询越多越好”。
- 变更流:把数据库中的变更以可恢复的事件流暴露给应用,但它不是消息队列,也不是天然的 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 对单个文档的写入具有原子性。假设订单文档内部有如下不变量:
如果订单行和订单总额都位于同一个文档中,那么修改订单行和总额可以由一次 updateOne 完成,其他读者不会看到只修改了一半的状态。
但如果订单和账户余额分别位于两个文档中:
orders/order-1001
accounts/customer-7
则一次普通的 updateOne 不能同时保证以下不变量:
此时有三个选择:
- 重新建模,把真正需要共同原子更新的数据放到一个文档;
- 接受最终一致性,并使用补偿或对账;
- 使用多文档事务。
事务是第三种选择,不是第一种选择的替代品。
1.3 MongoDB 与 JDBC、JPA 的边界
JDBC 面向关系数据库的表、行、列和 SQL;Jakarta Persistence(JPA)面向实体、关系映射、持久化上下文和 JPQL。MongoDB Java 驱动不实现 JDBC 的 Connection、PreparedStatement 或 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
MongoDatabase 和 MongoCollection 是轻量级句柄,可以按需获取。不要每个 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);
BigDecimal、Instant 和自定义值对象必须验证驱动版本中的编码支持及其 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("订单不能从当前状态支付");
}
这里的条件更新同时承担了乐观并发控制:
如果两个线程同时执行,最多一个线程能匹配 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);
}
每一步成立的原因是:
- 订单过滤器要求状态为
CREATED,避免重复支付; - 库存过滤器要求可用库存至少为
quantity,避免扣成负数; - 两个操作共享同一个
ClientSession; withTransaction负责事务开始、正常提交以及符合驱动规范的事务重试路径;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)
它不等价于分别拥有 status 和 createdAt 的独立索引。字段顺序必须由真实过滤、排序和范围条件决定。
一个常用推导方式是:
- 先放选择性高且需要等值匹配的字段;
- 再考虑排序字段;
- 范围条件通常会限制后续字段对排序或定位的帮助;
- 用
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 索引上线和恢复
索引创建可能影响资源使用。生产变更应至少验证:
- 当前数据是否满足唯一索引约束;
- 创建过程的 CPU、磁盘、复制延迟;
- 创建完成后
listIndexes是否包含预期索引; - 目标查询的执行计划和延迟;
- 删除旧索引前是否仍有其他查询依赖它。
不要只根据索引名称判断它被使用。先从数据库分析、应用慢查询日志和 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]
变更流不是轮询集合。应用读取的是服务器生成的事件,并通过游标持续等待新事件。
常见事件类型包括:
insertupdatereplacedeleteinvalidate
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 可能已经无法恢复,服务器会返回无法继续恢复的错误。此时处理器不能假装从断点继续,通常需要:
- 重新做一次全量同步;
- 记录新的同步边界;
- 再从该边界启动变更流;
- 通过幂等逻辑处理全量同步和增量事件之间的重叠。
如果数据库启用了集合级的变更前后镜像功能,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 中没有已提交事件
如果订单状态更新成功但库存更新因过滤条件失败,而代码没有抛出异常,事务可能错误提交。因此,检查 matchedCount 或 modifiedCount 不是形式动作,而是业务状态转换的一部分。
7.3 为什么还需要 outbox
如果在订单事务提交后直接调用消息系统:
提交订单事务成功
-> 调用消息系统失败
数据库已有支付订单,但下游没有收到事件。
如果先调用消息系统,再提交数据库:
消息发送成功
-> 数据库事务失败
下游又会收到一个不存在的支付事件。
outbox 把业务状态和待发送事件放进同一个 MongoDB 事务,使二者一起提交。下游发送器可以轮询 outbox,或通过 outbox 集合的变更流触发发送。发送器要使用唯一事件 ID 和幂等处理,避免重启时重复产生不可逆副作用。
八、错误表现与诊断路径
8.1 “事务不工作”
常见表现:
MongoTransactionException
Transactions are not supported
诊断顺序:
- 确认连接到的是副本集或分片集群;
- 查看服务端版本;
- 确认所有事务操作都使用同一个
ClientSession; - 确认没有把操作发到 secondary;
- 查看错误码、错误标签和服务端日志;
- 检查事务是否超过超时或涉及不支持的操作组合。
本地测试可以启动复制集后连接:
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 读取时可能出现解码异常,或者查询条件无法命中文档。修复方式不是在每个业务方法中添加隐式转换,而是:
- 统计字段类型分布;
- 编写迁移脚本统一 BSON 类型;
- 在写入入口强制类型;
- 增加数据库校验或集成测试。
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 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java Elasticsearch 工程:客户端、Mapping、查询、Bulk 和重试
- 下一篇:Spring 参数校验:Bean Validation、分组、嵌套、错误模型和国际化
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论