第24章:Mongo Change Streams 实战——订单变化实时推送

第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。每个事件包含operationTypeinsert/update/delete…、fullDocument变更后的完整文档、ns发生变更的集合等信息。小胖那这不是跟 Kafka 很像为什么不用 Kafka大师补集关系不是替代。Change Streams 负责数据库到应用的实时变更事件流——零代码侵入数据库操作自动产生事件。Kafka 负责应用到应用的异步消息——你可以在消费 Change Streams 事件后再投递到 Kafka 给更远的下游。两者组成事件驱动架构的管道。技术映射Change Streams 数据库的 CDCChange 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 可以回溯到断线前的位置。每个事件都有_idresumeToken客户端记录最后一次成功处理的 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}]}}]constchangeStreamdb.orders_cs.watch(pipeline,{fullDocument:updateLookup// update 事件中返回完整的更新后文档})print(Change Stream 已启动等待事件...)// 非阻塞方式取下一个事件mongosh 中 tryNext 不会阻塞leteventconstmaxEvents10letcount0functionprocessNext(){eventchangeStream.tryNext()if(event){countprint(\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 事件constinsertResultdb.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,maxRunSecs30){print(缓存同步 Worker 启动...)constpipeline[{$match:{operationType:{$in:[insert,update,replace]},fullDocument.status:{$exists:true}}}]conststreamdb.getCollection(collectionName).watch(pipeline,{fullDocument:updateLookup})letprocessedCount0letlastResumeTokennullconststartTimeDate.now()while(Date.now()-startTimemaxRunSecs*1000){consteventstream.tryNext()if(event){processedCountlastResumeTokenevent._id// 模拟同步到 RedisconstcacheKeyorder:${event.fullDocument.orderNo}print([Redis] SET${cacheKey}${event.fullDocument.status})// 模拟发送 App 推送if(event.operationTypeupdateevent.updateDescription?.updatedFields?.status){constnewStatusevent.updateDescription.updatedFields.statusprint([Push] 订单${event.fullDocument.orderNo}状态变更为:${newStatus})}}else{sleep(500)// 空闲时休眠 500ms}}stream.close()print(Worker 停止: 处理了${processedCount}个事件)return{processedCount,lastResumeToken}}// 启动 worker在 mongosh 中用非阻塞模式运行几秒constresultstartCacheSyncWorker(orders_cs,10)print(处理统计:,JSON.stringify(result))步骤四Resume Token——断线续听目标模拟断线后重新连接并从上次位置继续。// 1. 处理事件并记录 tokenletlastTokennullconststream1db.orders_cs.watch([],{fullDocument:updateLookup})// 处理前 3 个事件for(leti0;i3;i){leteventnull// 等待事件while(!event){eventstream1.tryNext()if(!event)sleep(100)}lastTokenevent._idprint(处理事件${i1}:,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 恢复监听...)conststream2db.orders_cs.watch([],{fullDocument:updateLookup,resumeAfter:lastToken// 从这个 token 之后开始})letrecovered0conststartDate.now()while(Date.now()-start5000){consteventstream2.tryNext()if(event){recoveredprint(恢复事件${recovered}:,event.operationType)if(recovered1)break}else{sleep(200)}}stream2.close()print(恢复处理了,recovered,个事件 (应 0))步骤五事件过滤——pipeline 的高级用法目标使用更多过滤条件精确监听事件。// 只在已完成状态下才触发事件constcompletedPipeline[{$match:{$or:[{operationType:insert,fullDocument.status:已完成},{operationType:update,updateDescription.updatedFields.status:已完成}]}}]constfilteredStreamdb.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:已完成}})constevfilteredStream.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 的事件流// 查看当前是不是分片集群constisShardeddb.runCommand({isdbgrid:1})print(是否是分片集群:,isSharded.ok1?是 (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模拟缓存同步 Workermongodb-lab/replicaset/ch24-resume-token.jsResume Token 断线续听mongodb-lab/replicaset/ch24-pipeline-filter.jsPipeline 事件过滤3.4 测试验证use local_life// 1. 验证 Change Stream 可用需要复制集环境try{consttestStreamdb.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 CDCDebezium代码侵入零数据库侧高每个写操作加代码零事件可靠性高基于 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 StreamsOplog 窗口要够大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: updateLookupMongoDB 在事件生成后额外做一次查询返回完整文档有微小的查询开销或者通过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 转 Decimal128Decimal128 不能直接从 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(ignoreUnknowntrue)② 新字段用新名称而非改名——不修改旧字段的类型只是新增对应的新字段如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 实战修炼与源码剖析