Angular PWA+AWS Lambda+Mongo Atlas的Change Stream监听异常求助
嘿,我来帮你搞定这两个头疼的问题,都是Serverless + MongoDB Change Streams里常见的坑,咱们一个个来捋:
问题1:本地用sls offline运行时收到重复通知
为啥会这样?
sls offline在本地模拟Lambda环境时,每次请求都可能重新加载你的模块(哪怕加了--skipCacheInvalidation,本地的缓存机制也不如云端稳定)。你的db.js会在mongoose连接成功后调用quote-watcher.js,每次调用都会新建一个Change Stream监听实例——多个实例同时盯着数据库,数据一变自然就会发多条重复通知。
怎么解决?
核心是确保Change Stream只被初始化一次,用单例模式来控制:
修改你的quote-watcher.js,加个全局变量标记是否已经创建过监听实例:
const { sendDataByTokens } = require('../bin/notify-devices'); // 用全局变量保存单例实例 let changeStreamInstance = null; module.exports = db => { 'use strict'; // 如果已经有实例了,直接返回,不再重复创建 if (changeStreamInstance) return; const taskCollection = db.collection('quotes'); changeStreamInstance = taskCollection.watch([], { fullDocument: 'updateLookup' }); changeStreamInstance.on('error', err => { console.log(err); // 出错后重置实例,允许后续重新创建 changeStreamInstance = null; }); changeStreamInstance.on('change', change => { console.log(change); // 你的原有处理逻辑... if (change.operationType === 'insert') { return sendDataByTokens( change.fullDocument.participants, change.operationType, change.fullDocument, change.ns.coll, null, { stage: change.fullDocument.currStage }, ); } if (change.operationType === 'update') { const { updatedFields } = change.updateDescription; return sendDataByTokens( change.fullDocument.participants, change.operationType, change.fullDocument, change.ns.coll, change.documentKey._id, { stage: change.fullDocument.currStage, updatedFields }, ); } }); };
另外,你可以调整sls offline的配置,让模块缓存更稳定:
在serverless.yml里添加:
plugins: - serverless-offline custom: serverless-offline: useChildProcesses: false
问题2:云端Lambda不触发Change Stream,偶尔返回旧数据
核心原因:Lambda的设计和Change Stream不兼容
Lambda是事件驱动的短暂执行环境——它只有在收到请求时才启动,处理完请求返回响应后,要么冻结执行环境,要么直接销毁。你设置了callbackWaitsForEmptyEventLoop: false,这会让Lambda在返回响应后立刻结束,根本不会等待Change Stream的异步监听任务。
偶尔能触发监听器,是因为AWS复用了之前冻结的执行环境,但这种复用完全不可靠(闲置一段时间就会被销毁),而且如果之前的Change Stream连接已经断了,就会读到旧数据。再加上市面上Lambda的超时时间最多15分钟(你设的是30秒),根本撑不住长期的Change Stream连接。
正确的解决方案:改用MongoDB Atlas Trigger
别再让Lambda主动监听了,换个更符合Serverless架构的方式:让MongoDB Atlas主动把变更事件推给Lambda。
在MongoDB Atlas控制台创建Change Stream Trigger:
- 进入你的集群,点击左侧菜单的Triggers,然后点Add Trigger
- 类型选Database,选择要监听的数据库和
quotes集合,勾选insert和update操作 - 目标选AWS Lambda,选择你已经部署的Lambda函数(需要先给Atlas配置AWS权限,控制台会引导你操作)
修改Lambda代码:
- 删掉
quote-watcher.js里的监听逻辑 - 新增一个Lambda处理函数,直接接收Atlas传过来的Change Stream事件,然后执行FCM推送逻辑:
// 比如在routes.js里加一个专门处理Atlas Trigger的路由,或者单独写一个Lambda handler exports.handleDbChange = async (event) => { const change = event; // Atlas会把完整的Change Stream对象传进来 console.log('Received DB change:', change); // 你的FCM推送逻辑... if (change.operationType === 'insert') { await sendDataByTokens( change.fullDocument.participants, change.operationType, change.fullDocument, change.ns.coll, null, { stage: change.fullDocument.currStage }, ); } if (change.operationType === 'update') { const { updatedFields } = change.updateDescription; await sendDataByTokens( change.fullDocument.participants, change.operationType, change.fullDocument, change.ns.coll, change.documentKey._id, { stage: change.fullDocument.currStage, updatedFields }, ); } return { statusCode: 200 }; };
- 删掉
备选方案(如果不想用Atlas Trigger):
用AWS ECS/EC2跑一个长期运行的服务来监听Change Stream,然后把变更事件发到AWS EventBridge,再由EventBridge触发Lambda。但这个方案需要额外维护服务器,不如Atlas Trigger省心。
注意:如果非要硬撑在Lambda里监听,你可以把callbackWaitsForEmptyEventLoop改成true,但这样Lambda会一直等待事件循环为空才结束,大概率会超时,而且会浪费很多资源和费用,绝对不推荐。
内容的提问来源于stack exchange,提问作者D.Kurapin

