如何让Node.js中的MongoDB变更流持续运行?
解决Mongoose变更流停止监听的问题
你的变更流偶尔静默停止,核心原因是未处理连接中断、错误事件,且缺少自动重连与断点恢复逻辑,同时长时间无变更时可能因连接闲置被服务器断开。以下是具体解决方案:
1. 监听变更流的错误与关闭事件
变更流遇到网络中断、服务器重启等情况时会触发error或close事件,若不处理这些事件,流会直接停止且无报错。需在事件触发时重新创建变更流。
2. 持久化Resume Token
你当前将resumeTokenRis存在内存中,应用重启或崩溃后会丢失,导致无法从上次中断的位置恢复监听。建议将其持久化到文件、Redis或MongoDB集合中。
3. 配置心跳与长连接参数
显式配置参数避免长时间无变更时连接被断开:
heartbeatFrequencyMS:设置心跳间隔(默认30000ms),缩短间隔确保服务器感知连接存活maxAwaitTimeMS:设置单次等待变更的最长时间(默认10000ms),到期后返回空结果维持连接活跃
4. 封装变更流创建逻辑
把创建变更流的代码封装成函数,方便在错误、关闭事件中调用重试。
修改后的完整代码示例
const fs = require('fs').promises; const RESUME_TOKEN_PATH = './resume-token-ris.json'; // 从文件读取持久化的resumeToken async function loadResumeToken() { try { const data = await fs.readFile(RESUME_TOKEN_PATH, 'utf8'); return JSON.parse(data); } catch (err) { // 文件不存在时返回null return null; } } // 保存resumeToken到文件 async function saveResumeToken(token) { try { await fs.writeFile(RESUME_TOKEN_PATH, JSON.stringify(token), 'utf8'); } catch (err) { console.error('Failed to save resume token:', err); } } // 创建变更流的函数 async function createChangeStream() { let resumeTokenRis = await loadResumeToken(); const ris = mongoose.connection.db.collection("ris"); try { const changeStreamRis = ris.watch([], { resumeAfter: resumeTokenRis, heartbeatFrequencyMS: 15000, // 15秒发送一次心跳 maxAwaitTimeMS: 30000 // 30秒无变更则返回空,保持连接 }); // 监听变更事件 changeStreamRis.on("change", async (change) => { if (change.operationType === "insert") { resumeTokenRis = change._id; await saveResumeToken(resumeTokenRis); // 持久化最新的token console.log("RI Inserted"); const ri = change.fullDocument; let userIds = ri.issueStat.taggedUser?.map(userId => userId.toString()) || []; userIds = [...new Set(userIds)].map(id => new ObjectId(id)); const userTokens = await NotificationToken.find({ userId: { $in: userIds } }); const user = await User.findById(ri.user); let tokens = userTokens.flatMap(userToken => userToken.tokens); if (tokens.length === 0) { console.log("No tokens found for this user."); return; } const message = { notification: { title: `New RI Assignment`, body: `Assigned by ${user.name} with ${ri.priority} priority.`, }, data: { id: `${ri.refId}`, type: "RI", }, tokens: tokens, }; try { const response = await admin.messaging().sendEachForMulticast(message); console.log(response); for (let index = 0; index < response.responses.length; index++) { const res = response.responses[index]; if (!res.success) { const failedToken = tokens[index]; console.log(`Failed token ID: ${failedToken}`); try { await NotificationToken.updateOne( { userId: { $in: userIds } }, { $pull: { tokens: failedToken } } ); console.log(`Token ${failedToken} deleted successfully`); } catch (error) { console.error(`Failed to delete token ${failedToken}: ${error}`); } } } } catch (error) { console.error("Failed to send notification:", error); } } }); // 监听错误事件,出错后重新创建流 changeStreamRis.on("error", async (err) => { console.error("Change stream error:", err); // 延迟几秒后重试,避免频繁重试 setTimeout(() => createChangeStream(), 5000); }); // 监听关闭事件,关闭后重新创建流 changeStreamRis.on("close", () => { console.log("Change stream closed, restarting..."); setTimeout(() => createChangeStream(), 5000); }); console.log("Change stream started successfully"); } catch (err) { console.error("Failed to create change stream:", err); setTimeout(() => createChangeStream(), 5000); } } // 启动变更流 createChangeStream();
额外注意事项
- 若使用MongoDB单节点,变更流的断点恢复功能受限,建议改用副本集或分片集群
- 生产环境中,可将resumeToken持久化到Redis或MongoDB集合,替代文件存储
- 重试逻辑可增加指数退避策略(如第一次等5秒,第二次10秒,第三次20秒),避免给数据库造成压力
内容的提问来源于stack exchange,提问作者Dugong98
相关产品推荐
相关产品推荐

