Change Streams 是 MongoDB 3.6 引入的功能,允许应用实时监听集合(或整个数据库)的变更事件。它是构建事件驱动架构、数据同步、缓存失效等场景的核心工具。
1. 基本原理
应用层 MongoDB
│ │
├─── watch("db.orders") ──────────→┤
│ (建立长连接) │
│ ←── 变更事件推送 ─────────────┤
│ │
│ 事件类型: insert / update / delete
│ 包含: 完整文档/变更差异 + Namespace + clusterTime
Change Streams 基于 oplog,因此要求 MongoDB 以复制集或分片集群方式运行。
2. 基础用法
2.1 Mongo Shell
const pipeline = [
{ $match: { "fullDocument.status": "paid" } }, // 过滤条件
{ $project: { "fullDocument": 1, "operationType": 1 } }
];
const changeStream = db.orders.watch(pipeline, {
fullDocument: "updateLookup" // update 事件获取完整文档
});
changeStream.on("change", (change) => {
console.log(change);
// {
// _id: { _data: "..." }, // Resume Token
// operationType: "insert",
// fullDocument: { ... },
// ns: { db: "shop", coll: "orders" },
// documentKey: { _id: ObjectId("...") }
// }
});
2.2 Java 驱动
MongoClient client = MongoClients.create("mongodb://localhost:27017");
MongoDatabase db = client.getDatabase("shop");
// 监听单个集合
MongoChangeStreamCursor<ChangeStreamDocument<Document>> cursor =
db.getCollection("orders").watch()
.fullDocument(FullDocument.UPDATE_LOOKUP)
.cursor();
while (cursor.hasNext()) {
ChangeStreamDocument<Document> change = cursor.next();
switch (change.getOperationType()) {
case INSERT -> handleInsert(change.getFullDocument());
case UPDATE -> handleUpdate(change.getFullDocument());
case DELETE -> handleDelete(change.getDocumentKey());
}
}
3. Resume Token 与容错
// 持久化 Resume Token,支持从断点恢复
BsonDocument resumeToken = null;
while (true) {
ChangeStreamIterable<Document> iterable;
if (resumeToken != null) {
iterable = db.getCollection("orders").watch().resumeAfter(resumeToken);
} else {
iterable = db.getCollection("orders").watch();
}
try (MongoChangeStreamCursor<ChangeStreamDocument<Document>> cursor = iterable.cursor()) {
while (cursor.hasNext()) {
ChangeStreamDocument<Document> change = cursor.next();
process(change);
resumeToken = change.getResumeToken(); // 保存最新 Token
saveTokenToRedis(resumeToken); // 持久化
}
} catch (MongoException e) {
// 连接断开,使用保存的 resumeToken 重连
log.warn("Stream interrupted, will retry with token: {}", resumeToken);
sleep(5000);
}
}
Resume Token 有效期:默认与 oplog 窗口相同(通常为 24-72 小时)。
4. 常见 CDC 架构
方案 1: MongoDB → Change Streams → 业务处理
用途: 缓存失效、审计日志、通知推送
方案 2: MongoDB → Change Streams → Kafka → 下游消费
用途: 数据仓库同步、全文检索(ES)
方案 3: MongoDB → Change Streams → 另一个 MongoDB
用途: 跨区域复制、数据迁移
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。