什么是 Change Stream

Change Stream 是 MongoDB 用于实现变更追踪的解决方案,类似于关系数据库的触发器,但原理不完全相同:


Change Stream 触发器
触发方式 异步 同步(事务保证)
触发位置 应用回调事件 数据库触发器
触发次数 每个订阅事件的客户端 1次(触发器)
故障恢复 从上次断点重新触发 事务回滚


MongoDB 从3.6版本开始支持了 Change Stream 能力(4.0、4.2 版本在能力上做了很多增强),用于订阅 MongoDB 内部的修改操作,change stream 可用于 MongoDB 之间的增量数据迁移、同步,也可以将 MongoDB 的增量订阅应用到其他的关联系统;比如电商场景里,MongoDB 里存储新的订单信息,业务需要根据新增的订单信息去通知库存管理系统发货。

使用 Change Stream 非常简单,mongo shell 封装了针对整个实例、DB、Collection 级别的订阅操作。

注意:集群环境一定要开启 enableMajorityReadConcern: true 配置文件中开启

  1. db.getMongo().watch() 订阅整个实例的修改
  2. db.watch() 订阅指定DB的修改
  3. db.collection.watch() 订阅指定Collection的修改

新建连接1发起订阅操作

  1. db.test.watch([], {maxAwaitTimeMS: 60000}) // 最多阻塞等待 1分钟

新建连接2写入新数据

  1. db.test.insert({x: 100})
  2. db.test.insert({x: 200})
  3. db.test.insert({x: 300})
  4. db.test.insert({x: 400})

上述 ChangeStream 结果里,_id 字段标识着 oplog 的某个位置,如果想从某个位置继续订阅,在 watch 时,通过 resumeAfter 指定即可。比如每个应用订阅了上述3条修改,但只有第一条已经成功消费了,下次订阅时指定第一条的 resume token 即可再次订阅到接下来的2条。

  1. db.test.watch([], {maxAwaitTimeMS: 60000,resumeAfter:{"_data" : "825E8DB54F000000012B022C0100296E5A1004BDB3105CB8B44265BB704E58E8EA8B0446645F696400645E8DB54FF368D1BC67ABE7AE0004" }})