第24章:Mongo Change Streams 实战——订单变化实时推送
1. 项目背景
业务场景:本地生活电商的用户抱怨——“我下完单不知道什么时候发货,每次都要手动刷新订单页面。”“优惠券快过期了也不提醒,白白浪费了。” 产品经理提出"订单状态变更实时推送"需求——用户下单后、支付成功、商家接单、骑手取货、订单完成,每个状态变更都要实时推送到 App。同时,数据部门要求订单状态变化后,同步刷新 Redis 缓存、更新 Elasticsearch 搜索索引、发送 Kafka 消息给风控系统。
开发的第一反应是——“在每个状态的更新接口里直接发消息不就行了?” 但这样会紧耦合:订单服务需要在代码里硬编码"改完状态后通知缓存、搜索、风控"等多个下游,任何一个下游挂掉都会影响订单状态的正常更新。用 Change Streams 可以解耦——订单服务只管改数据库,下游系统监听数据库的变更事件,各自治地响应。
痛点:不用 Change Streams 的下场——紧耦合的服务间调用,一点故障扩散到全局;状态变更通知遗漏或重复(消息没确认、网络抖动、服务重启);Kafka/Redis/ES 等下游有各自的消费节奏,耦合在一起很难做背压。
2. 项目设计
小胖(手机震个不停):大师!我把订单状态改完后直接在代码里调用 4 个下游服务——缓存刷新、搜索索引同步、消息推送、风控上报。结果风控服务一挂,整个订单状态更新接口全部超时!
大师:这是紧耦合的典型问题——订单服务不应该知道谁关心订单变化。正确的架构是——订单服务"喊一声"数据库变更,让关心的人自己去听。
小胖:“喊一声”?怎么喊?
大师:用 MongoDB 的 Change Streams。它基于复制集的 Oplog,将数据库的任何变更(Insert/Update/Delete/Replace)事件以实时流的形式暴露给应用层。你的订单服务只管写数据库,下游的缓存服务、搜索服务、消息推送各自 watch 订单集合的变更流,独立处理自己的逻辑。任何下游挂了不影响订单写入。
技术映射:Change Streams 是 MongoDB 的事件驱动架构基础设施——它把 Oplog 的变更事件转化为应用层可消费的流式 API。每个事件包含operationType(insert/update/delete…)、fullDocument(变更后的完整文档)、ns(发生变更的集合)等信息。
小胖:那这不是跟 Kafka 很像?为什么不用 Kafka?
大师:补集关系,不是替代。Change Streams 负责数据库到应用的实时变更事件流——零代码侵入,数据库操作自动产生事件。Kafka 负责应用到应用的异步消息——你可以在消费 Change Streams 事件后,再投递到 Kafka 给更远的下游。两者组成事件驱动架构的管道。
技术映射:Change Streams = 数据库的 CDC(Change Data Capture)。Kafka = 应用间消息队列。Change Streams → Kafka Connector(如 Debezium)是常见的大数据管道模式。
小白(追问):Change Streams 基于 Oplog,那它和直接读 Oplog 有什么不同?是不是更安全?
大师:直接读 Oplog 是内部实现细节——MongoDB 版本升级可能改变 Oplog 格式,你的代码就得跟着改。Change Streams 是官方稳定的 API——抽象掉了 Oplog 的内部结构,提供resumeToken做断线续听、fullDocument选项控制回调内容、pipeline做事件过滤。
小白:那如果监听的客户端断线了,重连之后会不会丢事件?
大师:不会——只要你的 Oplog 窗口足够大,Change Streams 可以回溯到断线前的位置。每个事件都有_id(resumeToken),客户端记录最后一次成功处理的 token,重连时从这个 token 之后继续。这就是At-Least-Once语义——事件可能重复但不会丢失,消费端必须做幂等处理。
大师(总结):Change Streams 三个要点——① 数据库变更自动变为事件流,解耦写入和下游;② resumeToken 实现断线续听不丢事件;③ 配合聚合管道可以过滤事件(只监听 status 变化而非所有字段变化)。
3. 项目实战
3.1 环境准备
Change Streams 需要复制集环境(单机不支持 Oplog,也就没有 Change Streams)。用第 17 章的 3 节点复制集。
dockercompose-fmongodb-lab/replicaset/docker-compose-rs.yml up-d3.2 分步实现
步骤一:启动 Change Streams 监听
目标:在 mongosh 中监听订单表的变更事件。
// 在 Primary 上连接的 mongoshuse local_life db.orders_cs.drop()// 开启 Change Streamconstpipeline=[{$match:{// 只监听 status 相关的 update$or:[{operationType:"insert"},{operationType:"update","updateDescription.updatedFields.status":{$exists:true}},{operationType:"replace"}]}}]constchangeStream=db.orders_cs.watch(pipeline,{fullDocument:"updateLookup"// update 事件中返回完整的更新后文档})print("Change Stream 已启动,等待事件...")// 非阻塞方式取下一个事件(mongosh 中 tryNext 不会阻塞)leteventconstmaxEvents=10letcount=0functionprocessNext(){event=changeStream.tryNext()if(event){count++print(`\n=== 事件${count}===`)print(" 操作类型:",event.operationType)print(" 集合:",event.ns.coll)print(" 文档ID:",event.documentKey?._id||event.documentKey)if(event.fullDocument){print(" 完整文档:",JSON.stringify(event.fullDocument).slice(0,150))}if(event.updateDescription){print(" 变更字段:",JSON.stringify(event.updateDescription.updatedFields))}print(" resumeToken:",event._id._data?.slice(0,30)+"...")}returnevent}// 初始检测processNext()步骤二:触发各种变更事件
目标:插入、更新、替换订单,观察 Change Stream 捕获到的事件。
// 在另一个 mongosh 窗口连接到 Primaryuse local_life// 1. 插入订单 → 触发 insert 事件constinsertResult=db.orders_cs.insertOne({orderNo:"CS_INSERT_001",userId:"U_001",status:"待支付",totalAmount:NumberDecimal("299.00"),createdAt:newDate()})print("插入:",insertResult.insertedId)// 2. 更新状态 → 触发 update 事件(仅当我们监听 status 时)db.orders_cs.updateOne({orderNo:"CS_INSERT_001"},{$set:{status:"已支付"},$currentDate:{paidAt:true}})print("更新状态: 待支付 → 已支付")// 3. 更新其他字段(不涉及 status)→ 不触发事件(被 pipeline 过滤了)db.orders_cs.updateOne({orderNo:"CS_INSERT_001"},{$set:{note:"这是备注"}})print("更新非status字段 → 应被过滤")// 回到 Change Stream 窗口执行 processNext() 查看捕获的事件while(processNext()){}步骤三:实现完整的事件消费循环(模拟后端服务)
目标:模拟一个消费者——订单状态变更后同步 Redis 缓存。
// === 模拟 Redis 缓存刷新消费者 ===// 这个模式展示了如何在实际后端代码中使用 Change StreamsfunctionstartCacheSyncWorker(collectionName,maxRunSecs=30){print("缓存同步 Worker 启动...")constpipeline=[{$match:{operationType:{$in:["insert","update","replace"]},"fullDocument.status":{$exists:true}}}]conststream=db.getCollection(collectionName).watch(pipeline,{fullDocument:"updateLookup"})letprocessedCount=0letlastResumeToken=nullconststartTime=Date.now()while(Date.now()-startTime<maxRunSecs*1000){constevent=stream.tryNext()if(event){processedCount++lastResumeToken=event._id// 模拟同步到 RedisconstcacheKey=`order:${event.fullDocument.orderNo}`print(`[Redis] SET${cacheKey}=${event.fullDocument.status}`)// 模拟发送 App 推送if(event.operationType==="update"&&event.updateDescription?.updatedFields?.status){constnewStatus=event.updateDescription.updatedFields.statusprint(`[Push] 订单${event.fullDocument.orderNo}状态变更为:${newStatus}`)}}else{sleep(500)// 空闲时休眠 500ms}}stream.close()print(`Worker 停止: 处理了${processedCount}个事件`)return{processedCount,lastResumeToken}}// 启动 worker(在 mongosh 中用非阻塞模式运行几秒)constresult=startCacheSyncWorker("orders_cs",10)print("处理统计:",JSON.stringify(result))步骤四:Resume Token——断线续听
目标:模拟断线后重新连接并从上次位置继续。
// 1. 处理事件并记录 tokenletlastToken=nullconststream1=db.orders_cs.watch([],{fullDocument:"updateLookup"})// 处理前 3 个事件for(leti=0;i<3;i++){letevent=null// 等待事件while(!event){event=stream1.tryNext()if(!event)sleep(100)}lastToken=event._idprint(`处理事件${i+1}:`,event.operationType,event.fullDocument?.orderNo||event.documentKey)}stream1.close()print("记录最后 token:",lastToken)// 触发更多事件db.orders_cs.insertOne({orderNo:"CS_AFTER_DISCONNECT",status:"待支付",createdAt:newDate()})// 2. 从 token 恢复(模拟重连)print("\n从 token 恢复监听...")conststream2=db.orders_cs.watch([],{fullDocument:"updateLookup",resumeAfter:lastToken// 从这个 token 之后开始})letrecovered=0conststart=Date.now()while(Date.now()-start<5000){constevent=stream2.tryNext()if(event){recovered++print(`恢复事件${recovered}:`,event.operationType)if(recovered>=1)break}else{sleep(200)}}stream2.close()print("恢复处理了",recovered,"个事件 (应 > 0)")步骤五:事件过滤——pipeline 的高级用法
目标:使用更多过滤条件精确监听事件。
// 只在"已完成"状态下才触发事件constcompletedPipeline=[{$match:{$or:[{operationType:"insert","fullDocument.status":"已完成"},{operationType:"update","updateDescription.updatedFields.status":"已完成"}]}}]constfilteredStream=db.orders_cs.watch(completedPipeline,{fullDocument:"updateLookup"})// 插入一条非"已完成"状态的 → 不被监听db.orders_cs.insertOne({orderNo:"CS_FILTERED_OUT",status:"待支付",createdAt:newDate()})// 更新到"已完成" → 触发db.orders_cs.updateOne({orderNo:"CS_FILTERED_OUT"},{$set:{status:"已完成"}})constev=filteredStream.tryNext()print("过滤后收到:",ev?`${ev.operationType}→${ev.fullDocument?.status}`:"无事件(可能还未到达)")filteredStream.close()步骤六:复制集 vs 分片集群中 Change Streams 的差异
// 1. 在复制集中:Change Stream 在 Primary 上启动// 事件从该节点的 Oplog 产生// 2. 在分片集群中:Change Stream 在 mongos 上启动// 事件从所有分片的 Oplog 合并产生// mongos 会协调来自多个 shard 的事件流// 查看当前是不是分片集群constisSharded=db.runCommand({isdbgrid:1})print("是否是分片集群:",isSharded.ok===1?"是 (mongos)":"否 (mongod)")// 在复制集中,Change Stream 连接到的 mongod 如果发生选举(Primary 切换)// 客户端需要重建 Change Stream —— Driver 会自动处理// 参数 resumeAfter / startAfter 可在重连时从断点继续3.3 完整代码清单
| 文件 | 用途 |
|---|---|
mongodb-lab/replicaset/ch24-change-stream.js | 基础 Change Stream 监听 |
mongodb-lab/replicaset/ch24-cache-sync.js | 模拟缓存同步 Worker |
mongodb-lab/replicaset/ch24-resume-token.js | Resume Token 断线续听 |
mongodb-lab/replicaset/ch24-pipeline-filter.js | Pipeline 事件过滤 |
3.4 测试验证
use local_life// 1. 验证 Change Stream 可用(需要复制集环境)try{consttestStream=db.orders_cs.watch([],{fullDocument:"updateLookup"})print("Change Stream 创建:","PASS")testStream.close()}catch(e){print("Change Stream 创建:","FAIL - ",e.message.includes("replica set")?"需要复制集":e.message)}// 2. 验证 insert 事件被捕获db.orders_cs.insertOne({orderNo:"VALIDATE_001",status:"测试",createdAt:newDate()})// 在 Change Stream 消费者中应看到此事件// 3. 验证 resumeToken 断线恢复// 手动记录一个 token,触发事件后,用 resumeAfter 重新连接// 应能看到断线期间的新事件// 4. 验证 pipeline 过滤// 插入 status != 目标值,应不被监听print("\n=== Change Streams 验证完成 ===")4. 项目总结
4.1 Change Streams vs 应用层 Hook vs Kafka CDC
| 维度 | Change Streams | 应用层 Hook(手动发事件) | Kafka CDC(Debezium) |
|---|---|---|---|
| 代码侵入 | 零(数据库侧) | 高(每个写操作加代码) | 零 |
| 事件可靠性 | 高(基于 Oplog) | 低(异常时易丢) | 高(Kafka 持久化) |
| 延迟 | < 100ms | < 1ms(同步) | < 500ms |
| 过滤能力 | 内置 pipeline | 应用层自由 | Debezium SMT |
| 运维复杂度 | 低(MongoDB 内置) | 无 | 高(需维护 Kafka + Connector) |
| 适用场景 | 实时推送、缓存同步 | 事件必须与写入原子绑定时 | 大数据管道、多系统分发 |
4.2 适用场景
Change Streams 适用:
- 订单/物流状态推送——状态变更实时通知用户 App。
- Redis 缓存刷新——数据变更后异步更新缓存。
- Elasticsearch 索引同步——MongoDB 主存 + ES 搜索引擎的 CDC 管道。
- 事件驱动架构的基础设施——其他服务通过 Change Streams 订阅领域事件。
- 实时数据看板——通过 Change Streams 将新数据流式推送至 Dashboard。
不适用场景:
- 需要严格的事务绑定(如"发消息必须在同一个数据库事务中成功")——Change Streams 是异步的、在事务提交后才触发。
- 一次性历史数据的全量导出——Change Streams 只关注增量变更。
4.3 注意事项
| 注意事项 | 说明 |
|---|---|
| 仅复制集/分片集群支持 | 单节点 mongod 没有 Oplog,不支持 Change Streams |
| Oplog 窗口要够大 | resumeToken只能回溯到 Oplog 窗口内,超出则 Change Stream 失效 |
| 消费端必须幂等 | Change Streams 提供 At-Least-Once 语义,网络重连可能发送重复事件 |
fullDocument: "updateLookup" | 如果文档在事件到达前被删除了,fullDocument 为 null |
| 不能 watch admin/local/config 库 | Change Streams 只支持业务数据库的集合 |
4.4 常见踩坑经验
故障案例一:resumeToken 失效后 Change Stream 静默丢失
某团队在 Worker 重启时用旧的 resumeToken 恢复监听,结果发现一批订单的状态变更没被推送到 App。根因:resumeToken 引用的 Oplog 条目已被覆盖——Oplog 窗口太小,重启后的几个小时恢复间隔超出了 Oplog 的保留时长。解决:增大 Oplog 至 24-48 小时;Worker 在每次处理完事件后持久化当前时间戳 + resumeToken,异常重启后检查"如果时间差过大,改用 startAtOperationTime 而非 resumeAfter"。
故障案例二:fullDocument: "default"导致 update 事件缺失完整文档
开发监听 update 事件,只拿到了updateDescription而不包含完整文档,无法判断订单所属用户来发送 App 推送。根因:fullDocument默认值是"default"——不返回完整文档。解决:显式设置fullDocument: "updateLookup",MongoDB 在事件生成后额外做一次查询返回完整文档(有微小的查询开销),或者通过documentKey._id在应用层补查。
故障案例三:分片集群中 change stream 的全局排序问题
某系统在分片集群中用 Change Streams 监听订单变更,期望所有事件按发生时间排序。实际上不保证跨分片的严格时序——shard A 的 10:00:02 的事件可能比 shard B 的 10:00:01 的事件更早到达 mongos。根因:分片集群中 Change Streams 从每个分片独立拉取事件后合并,合并的顺序是"事件到达 mongos 的时间"而非"事件在各自分片上的发生时间"。解决:依赖全局时间的业务逻辑不应依靠 Change Streams 的顺序保证;使用clusterTime自行排序和去重。
4.5 思考题
- 如果 Change Stream 消费者处理事件时抛出异常(如 Redis 不可用),应该立即重试、跳过还是停止消费?为什么?
- Change Stream 的
startAtOperationTime和resumeAfter有什么区别?在什么场景下用前者而不是后者?
(答案将在第 25 章末尾揭晓)
上一章思考题答案:
Double 转 Decimal128:Decimal128 不能直接从 Double 构造(
NumberDecimal(3.14)会把 Double 的浮点误差带进去),需先转为字符串再构造 Decimal128。迁移脚本:db.collection.aggregate([{$match:{field:{$type:"double"}}}, {$set:{field:{$convert:{input:{$toString:"$field"}, to:"decimal"}}}}, {$merge:{...}}]),使用$toString将 Double 转为字符串,再用$toDecimal转为 Decimal128——整个过程在聚合管道中原子完成。避免老服务因新字段报错:① 新字段用可选类型(Optional),老服务在反序列化时忽略未知字段(如 Jackson 的
@JsonIgnoreProperties(ignoreUnknown=true));② 新字段用新名称而非改名——不修改旧字段的类型,只是新增对应的新字段(如amount_v2: Decimal128同时保留amount: Double);③ 通过版本号灰度,路由层根据请求头或用户 ID 将流量分发到新旧服务。
延伸阅读与资源
MongoDB 实战进阶与内核修炼
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析
