08. MongoDB Change Streams 实时同步

MongoDB Change Streams: 实时监听数据变更、CDC 架构、Resume Token 与容错设计

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
  用途: 跨区域复制、数据迁移

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「mongodb」更多文章

  1. 11. MongoDB 安全认证与备份恢复
  2. 10. MongoDB 性能调优与运维监控
  3. 09. Spring Data MongoDB 实战